jayzhan211 opened a new pull request, #25724: URL: https://github.com/apache/datafusion/pull/25724
## Which issue does this PR close? - Part of #20773 and #24704. Does not close them. This is the first of two PRs split out of #25567, which carried the whole series and its measurements. It is self-contained: the option is off by default, so nothing changes for anyone who does not set it. A follow-up will extend the same machinery to single-stage aggregations. ## Rationale for this change A hash aggregation with millions of groups per partition spends most of its time finding group ids — 69-97% of aggregate compute on ClickBench Q15-18 and Q31-35. Past roughly two million groups per partition the cost per row falls off a cliff: on Q18 it goes from about 45 ns per row to 530 ns, because a dozen tables of hundreds of MB are growing at once. This PR lets a final aggregation stop growing one table. Once it reaches the configured number of groups it moves its state and all further input into 64 hash buckets and aggregates them one after another with a small, reused table. Each bucket is emitted, released, compacted, spilled and split again on its own. Measured on an Apple M4 Pro, 12 partitions, `hash_aggregate_bucket_threshold = 262144`, ClickBench `hits_partitioned`. Every query is run off/on back to back, because suite-level runs drift on this machine; 3 rounds of the fastest of 3 iterations. **ClickBench total 20.65 s → 17.25 s = 0.835x.** | faster | | slower | | |---|---|---|---| | Q32 | 0.28x | Q14 | 1.08x | | Q18 | 0.56x | Q5 / Q12 | 1.07x | | Q33 | 0.63x | Q1 | 1.06x | | Q15 / Q4 | 0.81x / 0.82x | Q10 / Q27 | 1.04x | | Q34 / Q8 | 0.87x | | | | Q35 / Q31 | 0.89x / 0.91x | | | | Q17 | 0.94x | | | **TPC-H SF10, same method: 5.02 s → 4.88 s = 0.972x.** Most of its aggregations are small or integer-keyed and never reach the threshold; the ones that move are Q17 0.81x and Q18 0.94x, against Q16 1.09x (a 47 ms query). Peak RSS is typically 0.4-0.9x, because a bucket is released as soon as it has been aggregated. Known costs, which is why the option is off by default: | case | time | peak RSS | |---|---|---| | `GROUP BY l_partkey, l_suppkey` at SF10 (the input repeats each group 7x) | 0.86x | **1.5x** | | ClickBench Q12 / Q14 (0.5-1M groups per partition, string keys) | 1.07-1.08x | 0.8-0.9x | Where the remaining cost is, per final input row on Q14: one table spends 78 ns finding group ids, bucketed it spends 26 ns, but 41 ns goes into moving the row to its bucket first (hash, gather, copy). So bucketing trades the cache miss for the move, and only wins once the table is big enough that the miss dominates — which is why Q18 (213 ns of group ids with one table) and Q32 (194 ns) win by 2-3x while a query at half a million groups per partition does not. On Q12 the CPU is in fact a wash — the final aggregate costs 76 ms more and the partial aggregate 68 ms less, summed over all partitions — and the 7% is wall clock: without buckets the final stage interns every row as it arrives, overlapped with the scan, while with buckets it only routes during the scan and aggregates after the input ends. For what it is worth, a smarter trigger does not fix this. A larger threshold, a self-timing trigger and holding the input back until it proves large were all built and measured; each only moved the loss to other queries. Engaging *later* is worse than never engaging at all: on Q33 the final aggregate takes 12.4 s of compute if it never buckets, 3.5 s if it buckets at 65k groups, and 25.8 s if it buckets at 1M, because a table that has already grown has already paid the penalty and then pays the routing on top. ## What changes are included in this PR? One commit per step: 1. `datafusion.execution.hash_aggregate_bucket_threshold` (groups, default 0 = off) and the bucketed final aggregation: `aggregates/final_buckets.rs` routes rows by hash with one gather per batch, compacts a bucket that has grown, and spills per bucket through `SpillManager`; `FinalHashAggregateStream` drives the bucket output. State with nested types is excluded, and a soft group limit disables it. 2. A partial aggregation table that reaches the same threshold is flushed downstream while its groups do not recur — about 1024 sampled group hashes per flush, kept for the last 64 flushes, and flushing stops once more than 20% of a sample comes back. 3. One hash table is reused across the buckets, and a bucket of unique groups is not compacted again. 4. `refactor:` the bucket driver moves into `aggregates/bucketed_aggregation.rs`. No behaviour change; it is what the follow-up PR builds on. 5. A final table whose input holds about one row per group moves into buckets at a quarter of the threshold, so little work is repeated. Buckets are still split again at the full threshold. 6. The table of a bucket keeps its group keys where its input already holds them. A bucket's rows sit in one contiguous block that is dropped as soon as that bucket has been aggregated, so its table stores a view into that block instead of copying every new key into a buffer of its own. Off for the table that compacts a bucket, whose purpose is to make it smaller, and off for a normal aggregation, whose table outlives every batch it has seen. This removes 17-19 ns per row on multi-column string keys: Q14's group ids 44.7 → 25.7 ns, Q16 54.8 → 46.9. The same change to the single-column string path (`ArrowBytesViewMap`) was measured and reverted: that path compares the stored bytes of a candidate on every probe, and the copy it would avoid is what packs those bytes together, so Q33 — URLs, which share a long prefix — lost 0.53x → 0.57x. The reason is recorded in `group_values/mod.rs`. ## What is the testing strategy for this PR? - Unit tests in `aggregates/final_buckets.rs` and `aggregates/hash_stream.rs`: bucketed results equal the single table for the final and partial streams, buckets split again, buckets spill under a memory limit, repeated groups are compacted, the partial flush stops when groups recur, and a bucket that spills mid-stream still returns every row once. - Unit tests in `group_values/multi_group_by/bytes_view.rs` that a borrowing builder still reads its values after the arrays they came from are dropped, and that `take_n` keeps every row readable. - `aggregate_bucketed.slt`: the same queries (FILTER, ROLLUP, ordered `array_agg`, DISTINCT, integer, `Utf8` and `Utf8View` keys) at 4 and 1 partitions, with the option off and at 100 groups, all sections identical, with `bucket_splits` and `table_flush_count` checked through `EXPLAIN ANALYZE`. - With the default temporarily set to 100 so every aggregation buckets, the aggregate unit tests, the `memory_limit` tests and the aggregate / group by / distinct sqllogictests pass. Two known exceptions: `aggregate_fuzz::streaming_aggregate_test` fails because it compares partial state rows one by one and a flushing partial aggregation legitimately emits a group more than once — it would have to merge partial rows first if the default ever became non-zero; and one tight-memory file (`ordered_aggregate_spill.slt`, limits of 500K-2M) fails intermittently under a loaded parallel run and passes in isolation. - `cargo fmt`, workspace clippy and the extended test suite pass. ## Are there any user-facing changes? A new experimental configuration option, `datafusion.execution.hash_aggregate_bucket_threshold`, off by default and documented in `configs.md` together with its known costs. New metrics on `AggregateExec`: `bucket_splits`, `bucket_compactions` and `table_flush_count`. -- 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]
