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]
