Skip to content

fix: evaluate small window partitions per batch in PartitionAggregateWindowExec - #8

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

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

Conversation

@msafonov

@msafonov msafonov commented Oct 9, 2026 •

Copy link
Copy Markdown
Collaborator

Problem

PartitionAggregateWindowExec (enabled by spark.comet.exec.window.partitionAggregate.enabled) handled every window partition through the buffered path when an expression is not constant within the partition, e.g.

FIRST_VALUE(x, TRUE) OVER (PARTITION BY id ORDER BY ts ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING)

With ~2.3 rows per partition (main_etl DimUserLoader, ~1.2e9 partitions) each partition grew and freed three memory reservations. In Comet each of those is a JNI call into Spark's TaskMemoryManager lock; async-profiler showed ~40% of the stage CPU there and ~45% in per-partition bookkeeping. On a 1/50 sample the stage took 1304 s vs 288 s with the operator disabled.

Change

  • Partitions that begin and end inside one input batch are evaluated together for all supported expressions except RANGE frames starting at an offset (N PRECEDING/N FOLLOWING): first/last/nth_value with and without IGNORE NULLS, sum/count/min/max/avg over suffix frames, ntile, percent_rank, cume_dist. No buffering or reservation is needed for them. Partitions crossing batch boundaries (and large ones) keep the buffered, spilling path unchanged.
  • Reservations grow in 1 MiB chunks and keep up to one chunk between partitions; unused reserved memory is returned to the pool before spilling, and spilling frees reservations as before.

Tests

  • New differential test against WindowAggExec: 6 partition-size layouts (all 1-row, 2-row, mixed small, small mixed with partitions spanning several batches, one huge partition) x batch sizes 1/3/7/64/1000 x with and without spilling, ~60 expressions incl. NULL values, NULL/tied ORDER BY keys, IGNORE NULLS, ROWS and RANGE frames.
  • Pool-call test with a counting MemoryPool: 10,500 small partitions produced 42,000 pool calls before; now bounded per batch. A large partition still spills.
  • Ignored micro-benchmark bench_small_partitions_suffix_frame (2M partitions, ~2.25 rows): 6.0 s / 8.0M pool calls before, 0.06 s / 1.3K pool calls after; WindowAggExec 2.0 s.
  • CometWindowExecSuite + CometPartitionAggregateWindowSuite: 142/142.

Cluster check (full DimUserLoader query, prod executor config, partitionAggregate enabled)

run app task-hours window stage wall
Comet joom-1.1-4 (last night) 25.0 h 61.1k s 1037 s
vanilla Spark 13.8 h 13.9k s 587 s
this PR 12.7 h 15.4k s 673 s

In the window stage the operator's own time is now ~0.3% of the profile, and JNI memory calls are ~0.1% (they were ~40%).
Compared with prod mart.dim_user (same inputs): 2,746,284,075 rows, same key set, and every column matches except first_device_id in ~0.026% of rows. All sampled mismatches (17,262) are ties in the query's ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY join_ts): both devices belong to the user and share the minimum join_ts, so the query is non-deterministic there and the difference is unrelated to this operator.

🤖 Generated with Claude Code

…WindowExec

With many tiny window partitions and a frame that is not constant within a
partition (for example FIRST_VALUE(x) IGNORE NULLS OVER (PARTITION BY id
ORDER BY ts ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING)), every partition
went through the buffered reverse pass and grew and freed three memory
reservations. Behind Comet's JNI memory pool each of those calls takes Spark's
memory manager lock; a stage with ~2.3 rows per partition ran 4.5x slower
than with DataFusion's WindowAggExec.

* Partitions that begin and end inside one input batch are now evaluated
  together for every supported expression except RANGE frames starting at an
  offset: first/last/nth_value (with and without IGNORE NULLS) and aggregates
  over suffix frames, ntile, percent_rank and cume_dist, without buffering or
  reserving the rows. Partitions crossing batch boundaries keep the buffered,
  spilling path.
* Reservations grow in 1 MiB chunks and keep up to one chunk between
  partitions; unused reserved memory is returned before spilling, and spilling
  frees the reservation as before.

Micro-benchmark (2M partitions, ~2.25 rows each): 6.0 s and 8.0M pool calls
before, 0.06 s and 1.3K pool calls after; WindowAggExec takes 2.0 s.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@github-actions github-actions Bot added the bug Something isn't working label Oct 9, 2026
@msafonov
msafonov merged commit d5779d9 into branch-1.1 Oct 9, 2026
37 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant