Repository navigation
perf: stream simple window functions over sorted input without per-partition slicing - #9
Merged
Merged
Conversation
…rtition slicing Add CometSortedWindowExec for windows whose expressions are all ROW_NUMBER, RANK, DENSE_RANK, or LEAD/LAG with a constant offset and default and without IGNORE NULLS. Partition and peer boundaries are computed per batch with vectorized comparisons, window columns are produced for the whole batch, and input columns pass through untouched, so tiny window partitions no longer pay for slicing every column, per-partition evaluator state and concatenation. Other windows keep BoundedWindowAggExec / PartitionAggregateWindowExec / WindowAggExec. spark.comet.exec.window.sorted.enabled (default true) switches it off. CometWindowExec now shows the native output rows and compute time. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What
CometSortedWindowExec, a streaming window operator for sorted input, used instead of DataFusion'sBoundedWindowAggExecwhen every window expression isROW_NUMBER,RANK,DENSE_RANK, orLEAD/LAGwith a constant offset (|offset| <= 1024) and constant default, withoutIGNORE NULLS. Everything else keepsBoundedWindowAggExec/PartitionAggregateWindowExec/WindowAggExec.Why
BoundedWindowAggExecslices every input column per window partition, keeps per-partition evaluator state keyed byVec<ScalarValue>, and concatenates the outputs. With 1-2 rows per partition (SCD lead over history, mongo dedup by key) that costs ~4 µs per partition, more with wide nested rows.How
Per input batch: partition and peer boundaries come from vectorized comparisons of adjacent rows (
distinct, or a comparator for nested types; the first row against the last key of the previous batch), the window columns are computed for the whole batch, input columns pass through untouched. Only counters, the last keys and, forLAG, the last values cross batches; a batch withLEADwaits until enough following rows arrived. Counters are produced in Spark's result type, so no cast projection is added. Buffered batches are accounted in the memory pool. Partition and order keys are the same normalized expressions the old path uses, so NULL/NaN/-0.0 equality is unchanged.spark.comet.exec.window.sorted.enabled(defaulttrue) switches it off.CometWindowExecnow shows nativeoutput_rows/elapsed_compute.Tests
BoundedWindowAggExec: partition sizes all-1, all-2, mixed, geometric, one 2500-row partition; batch splits 1, 2, 3, 5/1/2, 64, 4096, random; PARTITION BY one/two columns, a struct column, none; NULL keys and values, ORDER BY ties, ASC and DESC NULLS LAST; lead/lag offsets 0..3 with null/typed/cross-typed defaults on bigint/string/struct; empty batches and empty input; planner routing test.CometWindowExecSuite: prod shapes (lead with far-future timestamp default, row_number DESC NULLS LAST alone and under WindowGroupLimit + rn = 1, multi-column keys, rank/dense_rank ties, lead/lag 1..3, NaN/-0.0 and struct keys, nested lead values, IGNORE NULLS fallback), each with the operator on and off and batch sizes 3 and 8192.benches/sorted_window.rs(8 x 8192 rows, old -> new): lead, 1-row partitions 58 ms -> 0.43 ms (wide nested rows 283 ms -> 0.42 ms); row_number 49 ms -> 0.30 ms (wide 187 ms -> 0.25 ms); 2.2-row partitions lead 20.8 ms -> 0.35 ms, row_number 15.8 ms -> 0.23 ms.Cluster check (same input and config, joom-1.1-4 -> this branch)
sat_variant_warehouse_destinationsLEAD stage: 93.7 -> 36.6 s per task, 441M rows, result hashes equal.row_numberover wide nested rows, WindowGroupLimit + rn = 1): stage 226 -> 27 s, 14.3M rows, result hashes equal.🤖 Generated with Claude Code