andygrove opened a new pull request, #25843: URL: https://github.com/apache/datafusion/pull/25843
## Which issue does this PR close? - Part of #25758 - Related to #25157 ## Rationale for this change On `branch-55`, aggregating input that is sorted on a prefix of the group keys can be several times slower than on 54 (#25157, TPC-DS Q75). The new aggregation path, which is on by default in 55, emits completed groups `batch_size` at a time, and each emit shifts the state of all the remaining groups. The current workaround is `SET datafusion.execution.enable_migration_aggregate = false`. ## What changes are included in this PR? This PR backports #25312 from @2010YOUY01 to the `branch-55` line. `OrderedPartialAggregateStream` now materializes each run of completed groups once and emits slices of it. The cherry-pick conflicts in `common_ordered.rs` and `ordered_partial_table.rs`. On `main`, the new `materialize_groups` records per-accumulator timings through APIs that are not on `branch-55` (`AccumulatorPhase`, `aggregate_accumulator_metrics`, `time_emitting`). Here it takes `is_final: bool`, like the existing `next_output_batch_for_mode`, and records `emitting_time` as before. Otherwise the logic matches `main`. `ordered_partial_stream.rs` applied cleanly and matches `main` at #25312, apart from a doc comment that #25007 removed on `main`. This PR does not include #25639, which does the same for `OrderedFinalAggregateStream`. It builds on the #25538 spill refactor, so it would need a larger port, and most of the regression is recovered without it (see below). ## Are these changes tested? #25312 relies on existing tests. On this branch I ran: - `cargo test -p datafusion-physical-plan --lib` - `cargo test -p datafusion --test core_integration memory_limit` - `cargo test -p datafusion --features extended_tests --test fuzz -- aggregate`, which includes the sorted-input `streaming_aggregate_test` and the limited-memory aggregate fuzz tests - the full sqllogictest suite - `./dev/rust_lint.sh` I also timed the query from #25157 (10M rows, `SELECT count(*) FROM (SELECT DISTINCT ...)` over a Parquet file declared `WITH ORDER (d_year ASC)`), using `datafusion-cli` built with `--profile release-nonlto` on an M3 Max. Median of 5 runs: | Build | Time | |---|---| | `branch-55` | 2.91 s | | `branch-55`, `enable_migration_aggregate = false` | 0.31 s | | this PR | 0.43 s | ## Are there any user-facing changes? Ordered partial aggregation is faster. Output batches are still at most `batch_size` rows, except when the stream cannot reserve memory for a materialized run: then it passes the run downstream as one batch, as on `main`. No public API 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]
