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]
