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

   ### Describe the bug
   
   The native map lookup behind `element_at(map, key)` and `map[key]` can 
retain far more buffer capacity than its output needs when the map's values are 
nested, for example `map<string, array<int>>`. The values it returns are 
correct.
   
   This is the map counterpart of #6225. `spark_map_extract` gathers the 
matched values with a single Arrow `take` at the end 
([`map_extract.rs:227-231`](https://github.com/apache/datafusion-comet/blob/85316e843f1554cc36cd2f6ef0406f5d9a2dae84/native/spark-expr/src/map_funcs/map_extract.rs#L227-L231)).
 For list values, Arrow 59.3.0's `take_list` reserves child capacity as 
`child_data.len() / values.len() * indices.len()`. That is the average length 
of every value in the batch, not of the selected ones. So looking up a key 
whose array is short, while other keys hold long arrays, leaves the output with 
a large and mostly empty child buffer.
   
   #6498 fixes this for `list_extract` only, and its helpers are private to 
that file.
   
   ### Steps to reproduce
   
   On `main` at `85316e843f`, with a native test that calls `spark_map_extract` 
directly:
   
   1. Build a `MapArray` of type `map<int, array<int>>` with 8,192 rows.
   2. Give every row two entries. Key `1` maps to an empty array and key `2` 
maps to an array of 128 integers.
   3. Look up the scalar key `1`.
   4. Inspect `get_buffer_memory_size()` on the output and on its child array.
   
   The output holds 8,192 empty arrays and no child integers, but it retains 
2,129,924 bytes. 2,097,152 of those belong to the empty integer child. These 
are the same numbers #6225 reports for `list_extract`.
   
   If key `1` maps to a one-element array instead, the child holds 8,192 
integers (32 KiB) and still retains 2 MiB.
   
   ### Expected behavior
   
   Looking up short or empty nested values should not retain child buffers 
sized from the values that were not selected. Results, null handling and the 
lookup's current speed should stay the same.
   
   ### Additional context
   
   - The kernel came from #5806, which shipped in 1.1.0. I have not checked 
whether the 1.0.x path it replaced retained memory the same way.
   - Comet's native shuffle charges buffer capacity against its memory 
reservation, so the extra capacity can bring spills forward. It also shows up 
in the `native_allocated` metric.
   - Arrow 60.0.0 keeps the same `take_list` estimate, so the next Arrow 
upgrade will not fix this.
   - Two approaches already exist in the crate. #6498 compacts oversized 
buffers after `take`, and its helpers could move to a shared module and be 
called here. `take_extrema_values` in `array_extrema.rs` (#5403) reserves the 
exact selected lengths up front instead. A plain row-wise `MutableArrayData` 
gather is not a good alternative, because it measured 2 to 18 times slower than 
`take` on nested values.
   - A regression test can follow `test_list_extract_nested_retained_capacity` 
in #6498, which compares retained bytes against a row-wise gather.
   


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