ScalarValue has interior mutability but is intentionally used as hash key
(
flat_keys: &ArrayRef,
flat_values: &ArrayRef,
keys_offsets: &[i32],
values_offsets: &[i32],
keys_nulls: Option<&NullBuffer>,
values_nulls: Option<&NullBuffer>,
)
| 149 | |
| 150 | #[allow(clippy::allow_attributes, clippy::mutable_key_type)] // ScalarValue has interior mutability but is intentionally used as hash key |
| 151 | fn map_deduplicate_keys( |
| 152 | flat_keys: &ArrayRef, |
| 153 | flat_values: &ArrayRef, |
| 154 | keys_offsets: &[i32], |
| 155 | values_offsets: &[i32], |
| 156 | keys_nulls: Option<&NullBuffer>, |
| 157 | values_nulls: Option<&NullBuffer>, |
| 158 | ) -> Result<(ArrayRef, ArrayRef, OffsetBuffer<i32>)> { |
| 159 | let offsets_len = keys_offsets.len(); |
| 160 | let mut new_offsets = Vec::with_capacity(offsets_len); |
| 161 | |
| 162 | let mut cur_keys_offset = keys_offsets |
| 163 | .first() |
| 164 | .map(|offset| *offset as usize) |
| 165 | .unwrap_or(0); |
| 166 | let mut cur_values_offset = values_offsets |
| 167 | .first() |
| 168 | .map(|offset| *offset as usize) |
| 169 | .unwrap_or(0); |
| 170 | |
| 171 | let mut new_last_offset = 0; |
| 172 | new_offsets.push(new_last_offset); |
| 173 | |
| 174 | let mut keys_mask_builder = BooleanBuilder::new(); |
| 175 | let mut values_mask_builder = BooleanBuilder::new(); |
| 176 | for (row_idx, (next_keys_offset, next_values_offset)) in keys_offsets |
| 177 | .iter() |
| 178 | .zip(values_offsets.iter()) |
| 179 | .skip(1) |
| 180 | .enumerate() |
| 181 | { |
| 182 | let num_keys_entries = *next_keys_offset as usize - cur_keys_offset; |
| 183 | let num_values_entries = *next_values_offset as usize - cur_values_offset; |
| 184 | |
| 185 | let mut keys_mask_one = vec![false; num_keys_entries]; |
| 186 | let mut values_mask_one = vec![false; num_values_entries]; |
| 187 | |
| 188 | let key_is_valid = keys_nulls.is_none_or(|buf| buf.is_valid(row_idx)); |
| 189 | let value_is_valid = values_nulls.is_none_or(|buf| buf.is_valid(row_idx)); |
| 190 | |
| 191 | if key_is_valid && value_is_valid { |
| 192 | if num_keys_entries != num_values_entries { |
| 193 | return exec_err!( |
| 194 | "map_deduplicate_keys: keys and values lists in the same row must have equal lengths" |
| 195 | ); |
| 196 | } else if num_keys_entries != 0 { |
| 197 | let mut seen_keys = HashSet::new(); |
| 198 | |
| 199 | for cur_entry_idx in (0..num_keys_entries).rev() { |
| 200 | let key = ScalarValue::try_from_array( |
| 201 | &flat_keys, |
| 202 | cur_keys_offset + cur_entry_idx, |
| 203 | )? |
| 204 | .compacted(); |
| 205 | if seen_keys.contains(&key) { |
| 206 | // TODO: implement configuration and logic for spark.sql.mapKeyDedupPolicy=EXCEPTION (this is default spark-config) |
| 207 | // exec_err!("invalid argument: duplicate keys in map") |
| 208 | // https://github.com/apache/spark/blob/cf3a34e19dfcf70e2d679217ff1ba21302212472/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala#L4961 |
no test coverage detected
searching dependent graphs…