Skip to content

perf: stream simple window functions over sorted input without per-partition slicing - #9

Merged
msafonov merged 1 commit into
branch-1.1from
joom/1.1-fix-sorted-window
Oct 9, 2026
Merged

msafonov merged 1 commit into
branch-1.1from
joom/1.1-fix-sorted-window

Conversation

@msafonov

@msafonov msafonov commented Oct 9, 2026

Copy link
Copy Markdown
Collaborator

What

CometSortedWindowExec, a streaming window operator for sorted input, used instead of DataFusion's BoundedWindowAggExec when every window expression is ROW_NUMBER, RANK, DENSE_RANK, or LEAD/LAG with a constant offset (|offset| <= 1024) and constant default, without IGNORE NULLS. Everything else keeps BoundedWindowAggExec / PartitionAggregateWindowExec / WindowAggExec.

Why

BoundedWindowAggExec slices every input column per window partition, keeps per-partition evaluator state keyed by Vec<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, for LAG, the last values cross batches; a batch with LEAD waits 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 (default true) switches it off. CometWindowExec now shows native output_rows / elapsed_compute.

Tests

  • Rust differential test against 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.
  • Micro-benchmark 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_destinations LEAD stage: 93.7 -> 36.6 s per task, 441M rows, result hashes equal.
  • Mongo dedup (row_number over wide nested rows, WindowGroupLimit + rn = 1): stage 226 -> 27 s, 14.3M rows, result hashes equal.

🤖 Generated with Claude Code

…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>
@github-actions github-actions Bot added enhancement New feature or request performance labels Oct 9, 2026
@msafonov
msafonov merged commit 775c781 into branch-1.1 Oct 9, 2026
36 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant