jayzhan211 opened a new pull request, #25725: URL: https://github.com/apache/datafusion/pull/25725
## Which issue does this PR close? - Part of #20773 and #24704. Does not close them. **Stacked on #25724 — please review only the last commit.** The first six commits are that PR; GitHub shows them here because a fork PR cannot be based on another fork branch. This one is held as a draft until #25724 merges, at which point it will rebase down to a single commit. ## Rationale for this change #25724 lets a *final* hash aggregation move a table that has outgrown the CPU caches into hash buckets. The same problem exists for a single-stage aggregation — `AggregateMode::SinglePartitioned`, which is what a `GROUP BY` over already-partitioned input plans to, typically an aggregation over another aggregation. Its table grows the same way and pays the same cache miss per probe. This PR reuses the bucket driver extracted in #25724 (`aggregates/bucketed_aggregation.rs`) from `SingleHashAggregateStream`. A lone `Single` stream — one partition, no repartition below it — is deliberately **not** bucketed. Measured with `target_partitions = 1` it is 6-25% slower, because every row is then aggregated twice (once into the capped raw table, once into its bucket) and a single uncontended thread suffers far less from a large table in the first place. Evidence so far, on a nested aggregation over TPC-H SF10 `lineitem` (`GROUP BY l_orderkey, l_linenumber, mn` over a `GROUP BY l_orderkey, l_linenumber`), 12 partitions, threshold 262144: **0.85x time and 0.46x peak RSS**. Those numbers come from the whole series measured together; isolated before/after numbers for this commit alone will be added before this leaves draft. ## What changes are included in this PR? `SingleHashAggregateStream` gains the same three steps the final stream has: split its table into buckets once it reaches the threshold, route further input, and hand the buckets to `BucketedAggregation::output_stream` at end of input. It is enabled only when `agg.mode == AggregateMode::SinglePartitioned`, and the existing exclusions (nested state, soft group limit, no spill manager) carry over unchanged. ## What is the testing strategy for this PR? - Unit tests in `aggregates/hash_stream.rs`: a single-stage bucketed aggregation returns the same rows as one table, it spills its buckets under a memory limit, and a lone `Single` aggregation keeps its table. - `aggregate_bucketed.slt` gains a nested-aggregation section that asserts `mode=SinglePartitioned` with `bucket_splits`, run with the option off and at 100 groups with identical results, at both 4 and 1 partitions. ## Are there any user-facing changes? None beyond #25724: the same experimental option now also applies to single-stage aggregations that run on every partition. -- 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]
