jayzhan211 opened a new pull request, #25984:
URL: https://github.com/apache/datafusion/pull/25984

   ## Which issue does this PR close?
   
   - Related to #25157. No new issue: this is the Single-mode counterpart of 
the regression that #25639 fixed for `OrderedFinalAggregateStream` (and #25312 
for `OrderedPartialAggregateStream`).
   
   ## Rationale for this change
   
   With the default `enable_migration_aggregate = true`, an ordered 
single-stage aggregate is very slow when its sort prefix completes many groups 
at once. On the #25157 repro (10M rows sorted by `d_year`, 5 distinct `d_year` 
values, so ~2M groups complete at each boundary) with `target_partitions = 1`, 
the plan is
   
   ```text
   AggregateExec: mode=Single, gby=[], aggr=[count(1)]
     AggregateExec: mode=Single, gby=[d_year, brand, class, cat, manu, cnt, 
amt], aggr=[], ordering_mode=PartiallySorted([0])
   ```
   
   and the query takes 8.2 s on main and 0.55 s with this PR.
   
   `OrderedSingleAggregateStream` still took completed groups out of the live 
table `batch_size` rows at a time. Each `EmitTo::First(batch_size)` costs 
O(groups in the table), because group values are shifted and accumulators split 
off their tail. The stream also emitted at most one batch per input batch, so a 
large completed prefix was drained at EOF in O(N²/batch_size).
   
   ## What changes are included in this PR?
   
   - `OrderedSingleAggregateStream` materializes all completed groups once, 
then emits `batch_size` slices of that batch before it reads more input. A new 
`Outputting` state replaces `ProducingOutput`. This follows #25312 and #25639.
   - The materialized batch stays reserved together with the table until its 
last slice is handed off. If it cannot be reserved, it is handed off whole. 
Spilling is unchanged.
   - `OrderedAggregateTable<SingleMarker>::take_completed_result_batch` 
replaces `next_output_batch`. This PR removes the `batch_size` cap in 
`next_output_batch_inner` and the table's `batch_size` field, because nothing 
else uses them.
   
   ## What is the testing strategy for this PR?
   
   New unit tests in `ordered_single_stream.rs`, both failing on main:
   
   - `completed_groups_are_emitted_in_batch_size_slices`: with `batch_size = 
4`, one key boundary completes 43 groups. All of them are emitted before more 
input is read, and so are the groups completed at EOF. Every batch has at most 
`batch_size` rows, each group appears exactly once with the correct sum, and 
the reservation returns to 0.
   - `completed_groups_are_handed_off_whole_under_memory_pressure`: the memory 
limit is set one byte below the unlimited peak. The completed groups are then 
handed off as one batch without spilling, the results are correct, and the 
reservation returns to 0.
   
   Existing coverage passes: the `aggregates` unit tests, all sqllogictests, 
and the extended test suite.
   
   ### Benchmarks
   
   Apple M4 Pro. Both `datafusion-cli` binaries use the `release-nonlto` 
profile: one built from this branch, one from its base 801b017bb8. Each query 
file runs its query 3 times. The base and PR binaries were run alternately for 
5 rounds, and the table shows the median of 15 runs.
   
   | Case (`target_partitions`) | main | PR | PR / main |
   |---|---|---|---|
   | #25157 repro (1) | 8.170 s | 0.552 s | 0.07 |
   | #25157 repro (8), which does not use this stream | 0.410 s | 0.401 s | 
0.98 |
   | 20M rows sorted by `k = v / 4`, `GROUP BY k` (`Sorted`, ≤ `batch_size` 
groups completed per batch) (1) | 0.318 s | 0.318 s | 1.00 |
   | Same file, `GROUP BY k, x % 3` (`PartiallySorted([0])`) (1) | 0.573 s | 
0.562 s | 0.98 |
   | Peak RSS (`/usr/bin/time -l`), #25157 repro (1), median of 3 | 500 MB | 
434 MB | |
   
   The `ordered_group_values` criterion bench, group 
`fully_ordered_aggregate_exec` (Single mode, `Sorted` input), uses the same 
profile and 5 alternating rounds. It shows the median of the criterion point 
estimates.
   
   | Bench | main | PR | PR / main |
   |---|---|---|---|
   | int_run1 | 824.9 µs | 794.7 µs | 0.96 |
   | int_run8 | 402.8 µs | 405.9 µs | 1.01 |
   | int_run128 | 331.0 µs | 334.8 µs | 1.01 |
   | int_run8192 | 288.3 µs | 292.0 µs | 1.01 |
   | string_run1 | 1615.6 µs | 1591.4 µs | 0.99 |
   | string_run8 | 734.2 µs | 728.1 µs | 0.99 |
   | string_run128 | 582.7 µs | 585.1 µs | 1.00 |
   | string_run8192 | 526.6 µs | 524.0 µs | 1.00 |
   
   ## Are there any user-facing changes?
   
   Ordered single-stage aggregation is faster. Query results and public APIs 
are unchanged.
   


-- 
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