andygrove opened a new issue, #25851: URL: https://github.com/apache/datafusion/issues/25851
### Describe the bug A grouped aggregate whose groups carry large state, such as `array_agg` over a low-cardinality key, fails with `ResourcesExhausted` once it spills, even though the pool is free when it fails. With the same data spread over many small groups, the same query at the same memory limit spills and succeeds. The spill writes the aggregate's state in batches of at most `batch_size` rows, whatever their size in bytes ([stream.rs#L355-L359](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/datafusion/physical-plan/src/sorts/stream.rs#L355-L359), called from [spill.rs#L222-L230](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/datafusion/physical-plan/src/aggregates/spill.rs#L222-L230)). With fewer groups than `batch_size`, the whole table goes into one batch, and everything after the spill has to hold that batch at once: - With a single spill file, the merge reserves nothing ([multi_level_merge.rs#L329-L333](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/datafusion/physical-plan/src/sorts/multi_level_merge.rs#L329-L333)). The replay, `OrderedFinalAggregateStream` in `Sorted` mode, has no spill context, so when its table cannot grow to hold the batch, the query fails ([ordered_final_stream.rs#L332-L339](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/datafusion/physical-plan/src/aggregates/ordered_final_stream.rs#L332-L339)). This is what the repro below hits. - With two or more spill files, the merge first reserves `2 × max_record_batch_memory` per file ([multi_level_merge.rs#L540](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/datafusion/physical-plan/src/sorts/multi_level_merge.rs#L540), [sort.rs#L927-L934](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/datafusion/physical-plan/src/sorts/sort.rs#L927-L934)), and the same again as replay headroom ([L558](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/datafusion/physical-plan/src/sorts/multi_level_merge.rs#L558)). When not even one file fits, #25383 now tries to split it ([L596](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/datafusion/physical-plan/src/sorts/multi_level_merge.rs#L596)). But the split first reserves twice the batch it is about to split ([L682](https://github.com/apache/datafusion/blob/c1786f7e021fa2cc9b29fb2ffcbc0ec511302c55/d atafusion/physical-plan/src/sorts/multi_level_merge.rs#L682)), so a batch larger than half the budget still fails. In 55.1 this case returns the error directly ([55.1.0 multi_level_merge.rs#L525-L529](https://github.com/apache/datafusion/blob/55.1.0/datafusion/physical-plan/src/sorts/multi_level_merge.rs#L525-L529)). DataFusion Comet, which is on 55.1, hits this path. I have not reproduced it on main; the description here comes from reading the code. ### To Reproduce On main at c1786f7e0, save this as `repro.sql`: ```sql SET datafusion.execution.target_partitions = 2; SELECT count(*) AS groups, sum(n) AS values_collected FROM ( SELECT k, cardinality(array_agg(v)) AS n FROM ( SELECT value % 16 AS k, concat(CAST(value AS VARCHAR), repeat('x', 100)) AS v FROM generate_series(1, 3000000) ) GROUP BY k ); ``` `datafusion-cli -m 128M -f repro.sql` fails on every run: ``` Resources exhausted: Additional allocation failed for FinalHashAggregateStream[0] with top memory consumers (across reservations) as: AggregateStream[1]#9(can spill: false) consumed 48.0 B, peak 48.0 B, AggregateStream[0]#5(can spill: false) consumed 48.0 B, peak 48.0 B, AggregateStream[0]#2(can spill: false) consumed 48.0 B, peak 48.0 B. Error: Failed to allocate additional 214.5 MB for FinalHashAggregateStream[0] with 0.0 B already allocated for this reservation - 128.0 MB remain available for the total memory pool: greedy(used: 144.0 B, pool_size: 128.0 MB) ``` Controls with the same binary and the same data volume: - Without `-m`, the query returns 16 groups and 3,000,000 values. - With `value % 500000` (500,000 groups) and `-m 128M`, it returns 500,000 groups and 3,000,000 values. `EXPLAIN ANALYZE` shows the `FinalPartitioned` aggregate spilled 32 times (569.0 MB). To see where the memory goes, I added temporary `eprintln!` probes to `sort_and_spill`, the merge's admission, the final aggregate's out-of-memory branch and the replay loop. They are not needed for the repro. For each final partition, the output is (abbreviated): ``` final-oom: reserved=0 wanted=224915992 groups=8 input_rows=80 agg-spill: max_record_batch_memory=183944547 merge: spill_files=1 in_memory_streams=0 replay: input_rows=8 table_groups=8 resize_to=224914200 spillable=false ``` Here the partial aggregates hand each final partition an 80-row batch that already holds all of its state. So the final aggregate's first reservation fails, and it spills the whole table as one 8-row, 184 MB batch. Reading that back needs 225 MB in a 128 MB pool, although each group is only about 25 MB. ### Expected behavior The query succeeds under the memory limit, as it does when the same data is spread over many groups. A spilled run of a few large groups should be written in batches small enough to read back within the budget, for example by capping each spill batch by bytes as well as by rows. That can't help a single group that is larger than the budget, but here each group is about a fifth of the pool. ### Additional context - Found through DataFusion Comet, where `collect_list` and `collect_set` over a low-cardinality key fail this way after spilling while Spark completes the query. The Comet issue will link here. - Related: #25423 (replay table cannot spill; fixed in the test by capping merge fan-in), #25537 (one spill-replay driver) and #25428 (same merge code). - A variant of the repro looks like #25423's failure outside a test. Set `skip_partial_aggregation_probe_rows_threshold = 1` and `skip_partial_aggregation_probe_ratio_threshold = 0.0` in the repro. The spill batches are then small, and the merge and split succeed. But the query still fails, with the pool full (`used: 127.9 MB`): `Failed to allocate additional 10.0 MB for FinalHashAggregateStream[0] with 9.9 MB already allocated for this reservation - 75.1 KB remain available for the total memory pool`. I did not pin down that call site. -- 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]
