andygrove opened a new pull request, #25853:
URL: https://github.com/apache/datafusion/pull/25853

   ## Which issue does this PR close?
   
   - Closes #25851.
   
   ## Rationale for this change
   
   A grouped aggregate whose groups carry large state, such as `array_agg` over 
a low-cardinality key, fails with `ResourcesExhausted` after it spills, even 
though each group fits in memory and the same data spread over many small 
groups succeeds.
   
   Spilling writes the aggregate state in batches of at most `batch_size` rows, 
whatever their size in bytes. With fewer groups than `batch_size`, every group 
lands in one batch, and everything after the spill has to hold that batch at 
once:
   
   - With a single spill file, replay reads the batch directly and its table 
must hold every group at once. In the reproducer from the issue, that is 225 MB 
of state in a 128 MB pool, although each group is about 28 MB.
   - With several spill files, the merge reserves twice the largest batch of 
each file. When that does not fit, it tries to split the file, which first 
reserves twice the batch it is about to split.
   
   ## What changes are included in this PR?
   
   - Aggregate spill files are written in batches of at most `batch_size` rows 
and about 1 MiB (`SPILL_BATCH_TARGET_BYTES`). The number of rows per batch 
comes from the average row size of the spilled state, and is at least one. With 
the default `batch_size` of 8192, only rows averaging more than 128 bytes are 
written in smaller batches, so ordinary spills are unchanged.
   - The merge that replays the spill files uses the smallest number of rows 
per spilled batch as its batch size. Otherwise it would combine the rows of the 
small spilled batches back into one large batch before replay.
   
   This applies to every stream built on `AggregateSpill` (hash and ordered, 
final and single-stage). The legacy `GroupedHashAggregateStream`, used only 
when `datafusion.execution.enable_migration_aggregate` is `false`, is unchanged.
   
   ## What is the testing strategy for this PR?
   
   - New Case H in `aggregate_memory_spill.slt`: `array_agg` over 32 groups of 
about 2 MB each with a 24 MB limit. It runs once with two partitions (the final 
aggregation spills, and replay reads back one run per partition, as in the 
issue) and once with one partition (the single-stage aggregation spills five 
runs, and replay merges them). `EXPLAIN ANALYZE` assertions check that each 
spilled. Both queries fail on `main` with `ResourcesExhausted` and pass with 
this change. Reverting only the merge batch size makes the single-partition 
query fail again.
   - Unit test for how the rows per spilled batch are chosen: small rows keep 
`batch_size`, rows of a quarter of the target are written four to a batch, and 
a row larger than the target is written alone.
   
   With `datafusion-cli -m 128M`, the reproducer from the issue now returns 16 
groups and 3,000,000 values. It previously failed on every run. The variant 
from the issue's additional context, with partial aggregation skipped, also 
succeeds now.
   
   `external_aggr` shows no change on any of its 8 cases, since its rows are 
far below the target. On a spilling `array_agg` over 200,000 groups of about 
2.5 KB each, interleaved runs are 0.93x the time of `main` at 512 MB and 0.97x 
at 1 GB. At 512 MB, `main` also failed 7 of 15 runs with `ResourcesExhausted` 
and this change failed none. Spilling aggregates with narrow rows are unchanged 
(0.99x to 1.00x).
   
   ## Are there any user-facing changes?
   
   Aggregations whose groups have large state can now complete under memory 
limits where they previously failed after spilling. Spill files of rows 
averaging more than 128 bytes (with the default `batch_size`) contain more, 
smaller batches. No API or configuration changes.
   


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