sunchao opened a new issue, #6225:
URL: https://github.com/apache/datafusion-comet/issues/6225

   ### Describe the bug
   
   Follow-up to #5174.
   
   The Arrow `take` fast path in `list_extract` can retain substantially more 
buffer capacity when extracting short or empty inner arrays from an array whose 
other inner arrays are longer. The extracted values remain correct.
   
   This issue tracks that allocation behavior separately from the 
primitive/string extraction optimization in #5174.
   
   ### Steps to reproduce
   
   Compared PR head `78d65bd6a625b7616201daebc2f4bb5d0099264b` with base 
`3810936b40af3562003120742f7d258b234e4dc5`, using Arrow 59.3.0:
   
   1. Build an Arrow `ListArray` of type `array<array<int>>` with 8,192 rows.
   2. Each outer row contains two inner arrays: an empty array followed by an 
array of 128 integers. All rows and elements are valid.
   3. Extract the first inner array from every row with zero-based ordinal `0`, 
a null default, and `fail_on_error = false`. An absent default reaches the same 
`take` path.
   4. Compare the base `MutableArrayData` implementation with the PR's `take` 
implementation. Assert equal outputs, then inspect `get_buffer_memory_size()` 
on the output and its child array.
   
   Both implementations produce 8,192 valid empty arrays, with zero child 
integers.
   
   | Retained buffer capacity | Base | PR head |
   | --- | ---: | ---: |
   | Entire output | 65,600 bytes | 2,129,924 bytes |
   | Empty integer child | 32,768 bytes | 2,097,152 bytes |
   
   The input occupies 4,358,148 buffer bytes. The output capacity is 
approximately 32.5 times larger in the PR for this case.
   
   These measurements came from a native helper comparison using the exact 
base/head index and gather functions, with minimal scalar/error wrappers. It is 
a component test, not a Spark end-to-end memory benchmark.
   
   ### Expected behavior
   
   Selecting short or empty inner arrays should avoid retaining large child 
buffers sized from unselected data, while preserving the primitive/string 
extraction improvement.
   
   ### Additional context
   
   The new gather call is at 
[`list_extract.rs:392`](https://github.com/apache/datafusion-comet/blob/78d65bd6a625b7616201daebc2f4bb5d0099264b/native/spark-expr/src/array_funcs/list_extract.rs#L362-L392).
   
   In Arrow 59.3.0, `arrow_select::take::take_list` estimates child capacity as:
   
   ```text
   child_data.len() / values.len() * indices.len()
   ```
   
   Here that is `1,048,576 / 16,384 * 8,192 = 524,288` integers, reserving 2 
MiB even though the selected inner arrays contain no integers. The returned 
child array retains this allocation.
   
   Comet's native shuffle [recursively charges buffer 
capacities](https://github.com/apache/datafusion-comet/blob/78d65bd6a625b7616201daebc2f4bb5d0099264b/native/shuffle/src/partitioners/multi_partition.rs#L147-L169)
 against its memory reservation. Larger retained capacity can therefore advance 
spill thresholds. These numbers are allocated buffer capacity, **not measured 
resident RSS**. No end-to-end increase in spills or query runtime has been 
measured.
   
   Possible approaches include tighter allocation sizing in Arrow or a targeted 
Comet gathering strategy for affected nested types.
   
   Suggested validation for the follow-up:
   
   - Add a bounded regression for heterogeneous inner-array lengths, including 
nullable outer arrays.
   - Verify that retained capacity falls for small selections and that values, 
null handling, and ANSI errors are unchanged.
   - Check whether the same sizing behavior affects map results or structs 
containing lists.
   - Benchmark nested extraction alongside primitive/string cases to preserve 
the useful fast path.
   


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