Repository navigation
perf: materialize completed groups once in OrderedSingleAggregateStream - #25984
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #25984 +/- ##
==========================================
+ Coverage 82.62% 82.64% +0.01%
==========================================
Files 1147 1147
Lines 445320 445490 +170
Branches 445320 445490 +170
==========================================
+ Hits 367964 368157 +193
+ Misses 55013 54972 -41
- Partials 22343 22361 +18 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
2010YOUY01
left a comment
There was a problem hiding this comment.
The changes LGTM, thank you for the improvements.
One minor suggestion: I'd prefer those UTs to be removed, they're asserting the physical representation of output, and it adds maintenance overhead during refactors (e.g. it requires rework after we have blocked state layout).
A better way to prevent regression was adding some microbenchmark, and ensure to run it in future structural changes in aggregation.
|
Thanks @2010YOUY01 ! |
Which issue does this PR close?
OrderedFinalAggregateStream#25639 fixed forOrderedFinalAggregateStream(and perf: Avoid copying when materializing output in OrderedPartialAggregateStream #25312 forOrderedPartialAggregateStream).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 byd_year, 5 distinctd_yearvalues, so ~2M groups complete at each boundary) withtarget_partitions = 1, the plan isand the query takes 8.2 s on main and 0.55 s with this PR.
OrderedSingleAggregateStreamstill took completed groups out of the live tablebatch_sizerows at a time. EachEmitTo::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?
OrderedSingleAggregateStreammaterializes all completed groups once, then emitsbatch_sizeslices of that batch before it reads more input. A newOutputtingstate replacesProducingOutput. This follows perf: Avoid copying when materializing output in OrderedPartialAggregateStream #25312 and perf: Improve output materializing speed forOrderedFinalAggregateStream#25639.OrderedAggregateTable<SingleMarker>::take_completed_result_batchreplacesnext_output_batch. This PR removes thebatch_sizecap innext_output_batch_innerand the table'sbatch_sizefield, 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: withbatch_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 mostbatch_sizerows, 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
aggregatesunit tests, all sqllogictests, and the extended test suite.Benchmarks
Apple M4 Pro. Both
datafusion-clibinaries use therelease-nonltoprofile: one built from this branch, one from its base 801b017. 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.target_partitions)k = v / 4,GROUP BY k(Sorted, ≤batch_sizegroups completed per batch) (1)GROUP BY k, x % 3(PartiallySorted([0])) (1)/usr/bin/time -l), #25157 repro (1), median of 3The
ordered_group_valuescriterion bench, groupfully_ordered_aggregate_exec(Single mode,Sortedinput), uses the same profile and 5 alternating rounds. It shows the median of the criterion point estimates.Are there any user-facing changes?
Ordered single-stage aggregation is faster. Query results and public APIs are unchanged.