You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
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, called from 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). 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). 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, sort.rs#L927-L934), and the same again as replay headroom (L558). When not even one file fits, fix: leave memory for aggregate spill replay #25383 now tries to split it (L596). But the split first reserves twice the batch it is about to split (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). 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.
SETdatafusion.execution.target_partitions =2;
SELECTcount(*) AS groups, sum(n) AS values_collected FROM (
SELECT k, cardinality(array_agg(v)) AS n
FROM (
SELECT value % 16AS k, concat(CAST(value ASVARCHAR), 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):
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.
A variant of the repro looks like aggregate_memory_spill.slt Case G fails intermittently with ResourcesExhausted #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.
Describe the bug
A grouped aggregate whose groups carry large state, such as
array_aggover a low-cardinality key, fails withResourcesExhaustedonce 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_sizerows, whatever their size in bytes (stream.rs#L355-L359, called from spill.rs#L222-L230). With fewer groups thanbatch_size, the whole table goes into one batch, and everything after the spill has to hold that batch at once:OrderedFinalAggregateStreaminSortedmode, has no spill context, so when its table cannot grow to hold the batch, the query fails (ordered_final_stream.rs#L332-L339). This is what the repro below hits.2 × max_record_batch_memoryper file (multi_level_merge.rs#L540, sort.rs#L927-L934), and the same again as replay headroom (L558). When not even one file fits, fix: leave memory for aggregate spill replay #25383 now tries to split it (L596). But the split first reserves twice the batch it is about to split (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). 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 c1786f7, save this as
repro.sql:datafusion-cli -m 128M -f repro.sqlfails on every run:Controls with the same binary and the same data volume:
-m, the query returns 16 groups and 3,000,000 values.value % 500000(500,000 groups) and-m 128M, it returns 500,000 groups and 3,000,000 values.EXPLAIN ANALYZEshows theFinalPartitionedaggregate spilled 32 times (569.0 MB).To see where the memory goes, I added temporary
eprintln!probes tosort_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):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
collect_listandcollect_setover a low-cardinality key fail this way after spilling while Spark completes the query. The Comet issue will link here.aggregate_memory_spill.sltCase G fails intermittently with ResourcesExhausted #25423 (replay table cannot spill; fixed in the test by capping merge fan-in), Unify the spilling aggregate streams behind one spill-replay driver #25537 (one spill-replay driver) and perf: avoid unnecessary aggregate spill rewrites #25428 (same merge code).aggregate_memory_spill.sltCase G fails intermittently with ResourcesExhausted #25423's failure outside a test. Setskip_partial_aggregation_probe_rows_threshold = 1andskip_partial_aggregation_probe_ratio_threshold = 0.0in 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.