MCPcopy Create free account
hub / github.com/apache/datafusion / map_deduplicate_keys

Function map_deduplicate_keys

datafusion/spark/src/function/map/utils.rs:151–235  ·  view source on GitHub ↗

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>,
)

Source from the content-addressed store, hash-verified

149
150#[allow(clippy::allow_attributes, clippy::mutable_key_type)] // ScalarValue has interior mutability but is intentionally used as hash key
151fn 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

Callers 1

Calls 13

newFunction · 0.85
compactedMethod · 0.80
lenMethod · 0.45
mapMethod · 0.45
firstMethod · 0.45
pushMethod · 0.45
skipMethod · 0.45
iterMethod · 0.45
is_validMethod · 0.45
containsMethod · 0.45
insertMethod · 0.45
intoMethod · 0.45

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…