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]

Reply via email to