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]

Reply via email to