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]

Reply via email to