neilconway commented on code in PR #25607:
URL: https://github.com/apache/datafusion/pull/25607#discussion_r4094797602


##########
datafusion/functions-nested/src/array_compact.rs:
##########
@@ -112,101 +114,462 @@ fn array_compact_inner(arg: &[ArrayRef]) -> 
Result<ArrayRef> {
     let [input_array] = take_function_args("array_compact", arg)?;
 
     match &input_array.data_type() {
-        List(field) => {
-            let array = as_list_array(input_array)?;
-            compact_list::<i32>(array, field)
-        }
-        LargeList(field) => {
-            let array = as_large_list_array(input_array)?;
-            compact_list::<i64>(array, field)
-        }
+        List(field) => compact_list::<i32>(input_array, field),
+        LargeList(field) => compact_list::<i64>(input_array, field),
         Null => Ok(Arc::clone(input_array)),
         array_type => exec_err!("array_compact does not support type 
'{array_type}'."),
     }
 }
 
 /// Remove null elements from each row of a list array.
+///
+/// Each row is a range in a shared child array. Compaction removes null child
+/// values and rebuilds the row offsets, preserving null list rows.
 fn compact_list<O: OffsetSizeTrait>(
-    list_array: &GenericListArray<O>,
-    field: &Arc<arrow::datatypes::Field>,
+    input_array: &ArrayRef,
+    field: &FieldRef,
 ) -> Result<ArrayRef> {
+    let list_array = as_generic_list_array::<O>(input_array.as_ref())?;
+    let visible_len = offset_span_len(list_array.offsets());
+    if visible_len == 0 || list_array.null_count() == list_array.len() {
+        return Ok(Arc::clone(input_array));
+    }
+    // Restrict the child to the visible values before computing logical nulls,
+    // which can be expensive.
     let values = list_array.values();
+    let start = list_array.offsets()[0].as_usize();
+    let sliced_values = (start != 0 || visible_len != values.len())
+        .then(|| values.slice(start, visible_len));
+    let values = sliced_values.as_ref().unwrap_or(values);
     // Use logical nulls so element types without a validity buffer
     // (e.g. NullArray) are still treated as null.
-    let Some(values_nulls) = values.logical_nulls() else {
-        // Fast path: no validity buffer, no nulls to remove
-        return Ok(Arc::new(list_array.clone()));
+    let Some(values_nulls) = values
+        .logical_nulls()
+        .filter(|nulls| nulls.null_count() != 0)
+    else {
+        return Ok(Arc::clone(input_array));
     };
-    if values_nulls.null_count() == 0 {
-        // Fast path: validity buffer present but no nulls set
-        return Ok(Arc::new(list_array.clone()));
+    if values_nulls.null_count() == values_nulls.len() {
+        return Ok(empty_list_values(list_array, Arc::clone(field)));
     }
 
-    let list_nulls = list_array.nulls();
+    compact_list_values(list_array, field, values.as_ref(), &values_nulls)
+}
+
+/// Copy primitive values and build list offsets in one pass. For other types,
+/// build offsets and a keep bitmap, then copy Utf8/LargeUtf8 strings in valid
+/// spans or use Arrow's filter kernel for the remaining types.
+fn compact_list_values<O: OffsetSizeTrait>(
+    list_array: &GenericListArray<O>,
+    field: &FieldRef,
+    values: &dyn Array,
+    values_nulls: &NullBuffer,
+) -> Result<ArrayRef> {
+    let (offsets, new_values) = downcast_primitive_array! {
+        values => Ok(compact_primitive(values, list_array, values_nulls)),
+        _ => compact_non_primitive(values, list_array, values_nulls),
+    }?;
+    Ok(Arc::new(GenericListArray::<O>::try_new(
+        Arc::clone(field),
+        offsets,
+        new_values,
+        list_array.nulls().cloned(),
+    )?))
+}
+
+fn compact_primitive<T: ArrowPrimitiveType, O: OffsetSizeTrait>(
+    values: &PrimitiveArray<T>,
+    list_array: &GenericListArray<O>,
+    values_nulls: &NullBuffer,
+) -> (OffsetBuffer<O>, ArrayRef) {
     let list_offsets = list_array.offsets();
-    let original_data = values.to_data();
-    let (first_offset, visible_len) = offset_span(list_offsets);
-    let capacity =
-        visible_len - values_nulls.slice(first_offset, 
visible_len).null_count();
-    let mut offsets = Vec::<O>::with_capacity(list_array.len() + 1);
+    let first_offset = list_offsets[0].as_usize();
+    let mut output = Vec::with_capacity(values_nulls.len() - 
values_nulls.null_count());
+    let mut offsets = Vec::with_capacity(list_array.len() + 1);
     offsets.push(O::zero());
-    let mut mutable = MutableArrayData::with_capacities(
-        vec![&original_data],
-        false,
-        Capacities::Array(capacity),
-    );
-
-    for row_index in 0..list_array.len() {
-        if list_nulls.is_some_and(|n| n.is_null(row_index)) {
-            offsets.push(offsets[row_index]);
-            continue;
+    for (row, window) in list_offsets.windows(2).enumerate() {
+        if list_array.is_valid(row) {
+            let start = window[0].as_usize() - first_offset;
+            let end = window[1].as_usize() - first_offset;
+            for i in start..end {
+                if values_nulls.is_valid(i) {
+                    output.push(values.value(i));
+                }
+            }
         }
+        offsets.push(O::usize_as(output.len()));
+    }
+    let output = PrimitiveArray::<T>::new(output.into(), None)
+        .with_data_type(values.data_type().clone());
+    (OffsetBuffer::new(offsets.into()), Arc::new(output))
+}
 
-        let start = list_offsets[row_index].as_usize();
-        let end = list_offsets[row_index + 1].as_usize();
-        let row_null_count = values_nulls.slice(start, end - 
start).null_count();
-        let kept = (end - start) - row_null_count;
-
-        // Batch consecutive non-null elements into single extend() calls
-        // to reduce per-element overhead. For [1, 2, NULL, 3, 4] this
-        // produces 2 extend calls (0..2, 3..5) instead of 4 individual ones.
-        let mut batch_start: Option<usize> = None;
-        for i in start..end {
-            if values_nulls.is_null(i) {
-                // Null breaks the current batch — flush it
-                if let Some(bs) = batch_start {
-                    mutable.try_extend(0, bs, i)?;
-                    batch_start = None;
+#[inline(never)]
+fn compact_non_primitive<O: OffsetSizeTrait>(
+    values: &dyn Array,
+    list_array: &GenericListArray<O>,
+    values_nulls: &NullBuffer,
+) -> Result<(OffsetBuffer<O>, ArrayRef)> {
+    let first_offset = list_array.offsets()[0].as_usize();
+    let mut offsets = Vec::with_capacity(list_array.len() + 1);
+    offsets.push(O::zero());
+    let mut count = 0;
+    // Child validity is already the keep mask unless null parents hide values.
+    let mut mask =
+        (list_array.null_count() != 0).then(|| 
BooleanBufferBuilder::new(values.len()));
+    for (row, window) in list_array.offsets().windows(2).enumerate() {
+        let start = window[0].as_usize() - first_offset;
+        let len = window[1].as_usize() - window[0].as_usize();
+        if list_array.is_valid(row) {
+            let bit_start = values_nulls.offset() + start;
+            count += values_nulls.buffer().count_set_bits_offset(bit_start, 
len);
+            if let Some(mask) = &mut mask {
+                // Fill the gap left by any preceding null parent rows.
+                if start > mask.len() {
+                    mask.append_n(start - mask.len(), false);
                 }
-            } else if batch_start.is_none() {
-                batch_start = Some(i);
+                mask.append_packed_range(
+                    bit_start..bit_start + len,
+                    values_nulls.validity(),
+                );
             }
         }
-        // Flush any remaining batch after the loop
-        if let Some(bs) = batch_start {
-            mutable.try_extend(0, bs, end)?;
-        }
-
-        offsets.push(offsets[row_index] + O::usize_as(kept));
+        offsets.push(O::usize_as(count));
     }
+    let mask = mask
+        .map(|mut mask| {
+            mask.append_n(values.len() - mask.len(), false);
+            mask.finish()
+        })
+        .unwrap_or_else(|| values_nulls.inner().clone());
+    let output = match values.data_type() {
+        DataType::Utf8 => copy_string_spans(values.as_string::<i32>(), &mask, 
count),
+        DataType::LargeUtf8 => copy_string_spans(values.as_string::<i64>(), 
&mask, count),
+        _ => filter(values, &BooleanArray::new(mask, None))?,

Review Comment:
   I filed https://github.com/apache/arrow-rs/issues/11200 for this upstream



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to