Skip to content

fix: native windows reserve the rows they buffer and buffer one window partition at a time - #6279

Open
andygrove wants to merge 4 commits into
apache:mainfrom
andygrove:fix-6253-window-memory
Open

andygrove wants to merge 4 commits into
apache:mainfrom
andygrove:fix-6253-window-memory

Conversation

@andygrove

@andygrove andygrove commented Sep 27, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6253.

Rationale for this change

PERCENT_RANK, CUME_DIST, NTILE, and aggregates whose frame ends at UNBOUNDED FOLLOWING run in DataFusion's WindowAggExec. That includes Spark's default frame for agg(x) OVER (PARTITION BY k). WindowAggExec buffers the task's entire input, concatenates it, and reserves nothing. A large or skewed task therefore grows native memory that neither Comet's pool nor Spark can see.

Reserving those batches is not enough on its own. WindowAggExec holds the whole task input, not one window partition, so a reservation would fail any task whose input is bigger than its pool share, even when every window partition is small. A build that only added the reservation failed SUM(v) OVER (PARTITION BY k) over 200,000 rows in 5-row window partitions under a 4 MB pool.

What changes are included in this PR?

  • The planner now uses a new CometWindowAggExec where it used WindowAggExec. It evaluates the same WindowExprs one window partition at a time, as WindowAggExec does, with two differences:
    • It emits window partitions as they complete. The input is sorted by the PARTITION BY keys, so once a batch starts a new window partition, every earlier one is complete. Those are evaluated and emitted at once, and only the window partition that may continue is kept. Without PARTITION BY, it still buffers the whole input.
    • It reserves what it buffers. The buffered batches are counted with RecordBatchMemoryCounter, so buffers shared between batches are counted once. The copy made to concatenate them is reserved as soon as it is made. It counts only the buffers the copy does not share with the batches it was made from, which is exact for nested, dictionary and view types. Estimating each batch's share instead would charge every batch sliced from one list array for all of that array's values. Both are released once the output is emitted.
  • A refused reservation now says that the operator cannot spill and which settings to change.
  • BoundedWindowAggExec is unchanged. It keeps only the rows its frames still need, which is bounded by the frame rather than the window partition (for a RANGE ... CURRENT ROW frame, that is a whole group of peer rows). It still reserves nothing.
  • The operator tuning guide gets a short section on window functions.

A window partition that does not fit now fails the task with a memory error instead of growing untracked. Spark would spill in that case. Spilling for WindowAggExec is still apache/datafusion#22946. In on-heap mode, which the Spark SQL tests use, the pool is unbounded, so nothing is ever refused there.

Measured in release mode against WindowAggExec with a standalone bench over 2M rows in 8192-row batches:

  • Speed. Within 5% for 10-row and 100,000-row window partitions and without PARTITION BY. With 1,000-row window partitions it is 6–15% slower, about 0.5 ms per 2M rows, because it concatenates per batch rather than once.
  • Memory. With 20M rows (480 MB) in 1,000-row window partitions:
Peak RSS Peak reserved Largest output batch
WindowAggExec 1,254 MB 0 20M rows
CometWindowAggExec 480 MB (the input itself) 609 KB 9,000 rows

Without PARTITION BY, both peak at 1,090 MB, but CometWindowAggExec reserves 960 MB of it.

How are these changes tested?

  • Rust unit tests in window_agg.rs.

    • A differential test against WindowAggExec over 20 seeds. It uses random window partition sizes, random batch boundaries including empty batches, null keys and two key columns, with and without PARTITION BY.
    • Tests for emission order and for empty input.
    • Memory tests:
      • many small window partitions pass under a pool 50 times smaller than the input;
      • a window partition that does not fit is refused, with the new message;
      • the concatenated copy is reserved;
      • batches sliced from one list array are charged once for the list values they share. A per-batch get_slice_memory_size estimate over-charged this case about tenfold and was refused.
    • Each of five mutations fails at least one of these tests: no batch reservation, no copy reservation, no release after emitting, never continuing a window partition across batches, and not merging a continued window partition.
  • CometWindowExecSuite. The new tests:

    • Whole-partition functions over 7-row batches.
    • Many small window partitions under a 4 MB pool pass and match Spark. This test fails with a build that only adds the reservation.
    • A window partition that does not fit fails the task with Additional allocation failed for CometWindowAggExec.

    The suite passes on Spark 3.4, 3.5 and 4.1.

  • Other suites.

    • The 20 window SQL file tests pass.
    • TPC-DS q12, q20, q47, q53, q57, q63, q89 and q98, and their v2.7 variants, match the golden files at SF1 under all three join configurations.
    • Spark 4.1.3's DataFrameWindowFunctionsSuite and DataFrameWindowFramesSuite ran from Comet's classpath with Comet enabled on-heap. They gave the same result on main and on this branch: 74 pass, and the same 4 fail. Those 4 are the tests the Spark diff marks IgnoreComet or adjusts for CometWindowExec.

…w partition at a time

PERCENT_RANK, CUME_DIST, NTILE, and aggregates whose frame ends at UNBOUNDED FOLLOWING ran in
DataFusion's WindowAggExec, which buffers the task's whole input, concatenates it, and reserves
none of it.

Replace it with CometWindowAggExec. It evaluates the same window expressions, but emits each
window partition once the sorted input moves past it, and reserves the rows it buffers and the
copy it concatenates for evaluation. A window partition that does not fit now fails the task with
a memory error instead of growing untracked.
@github-actions github-actions Bot added the bug Something isn't working label Sep 27, 2026
@andygrove
andygrove marked this pull request as ready for review September 27, 2026 17:05

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: WindowAggExec retained the entire task input without memory reservations.
  • Design approach: CometWindowAggExec emits completed window partitions and reserves retained buffers plus concatenation copies.
  • Correctness / compatibility analysis: Found one reproducible P2 issue: nested-column slices can cause false out-of-memory failures. Partition processing otherwise follows the checked Spark window implementations across supported version families.
  • Key design decisions: Reusing DataFusion’s plan metadata and expression evaluators limits duplication. Spilling remains unsupported and is documented.
  • Implementation sketch: The full five-file diff adds the operator, changes planner selection, adds Rust and Scala tests, and updates tuning guidance.
  • Behavioral changes worth calling out: Completed partitions return earlier. Oversized reservations now fail the task. Additional partition scans and boundary-row copies add work per batch. BoundedWindowAggExec remains unchanged.
  • Suggested improvements: Make concatenation estimates respect nested slice offsets and add the regression described below.

Reviewed full SHA 0776a2451f6d2c6ca438c012d77af78242525acb against 605051ad239ef704f5f25d67910a446a6b6d7c70. The PR is non-draft. Snapshot and live discussion checks found no existing reviews or comments. Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-expression-pr, plus the operator contributor guide.

Exact-head CI: 17 checks succeeded, 13 were skipped, and six remained running. Rust tests and the native build passed. Comet Spark 4.1 suites and TPC-H/TPC-DS verification were running. Spark SQL jobs were skipped. No failures were reported.

Validation: All six added Rust tests passed in an isolated harness importing the unchanged operator with DataFusion 55.1.0 and Arrow 59.3.0. An additional differential reproduction confirmed the finding. Full Comet/JVM and Spark SQL suites were not run locally because of shared-filesystem space exhaustion. Author-reported integration results and performance timings were not independently reproduced.

let mut copy_size = 0;
for batch in &batches {
for column in batch.columns() {
copy_size += column.to_data().get_slice_memory_size()?;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Could the copy estimate account for nested slice offsets before summing get_slice_memory_size()? For a sliced ListArray, this method recursively counts the entire child array, including values outside the slice. Multiple batches sharing that child therefore charge it repeatedly, although concatenation copies only the referenced values. A reproduced COUNT(k) OVER () with a list payload fails under a 16 MiB pool even though its retained input and actual concatenated copy fit together. This introduces task failures based on batch representation rather than required memory. Please use offset-aware nested sizing and add a shared-list-slice regression.

Evidence: Executed an isolated harness against the unchanged exact-head operator. Constructed 10,000 rows with k = 0..9999 and an a: List<Int64> containing 0..63 per row, then supplied ten original.slice(i * 1000, 1000) batches to COUNT(k) OVER () with GreedyMemoryPool::new(16 * 1024 * 1024). Unique retained input measured 8,508,612 bytes and the actual concatenated batch measured 5,265,536 bytes, totaling less than 16 MiB. The new estimate requested 51,320,040 bytes and failed with Resources exhausted: Failed to allocate additional 48.9 MB for CometWindowAggExec[0]. The previous WindowAggExec succeeded. Materializing the same slices independently with take_record_batch also made the new operator succeed and return identical results.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, good catch. The copy is now reserved right after concatenating, counting only the buffers it doesn't share with the input batches, which is exact for nested types instead of estimated per slice. I added batches_sliced_from_one_list_array_are_charged_once with your repro's shape; it's refused with the old estimate and passes now.

evaluate() no longer re-evaluates the partition keys and boundaries that push_batch already
found, since the buffered rows are always one window partition, and the cross-batch check
compares only the key columns. Also drop the recomputed ordered_partition_by_indices and the
input_schema field, share one comparison helper in the Rust tests, and fold the tiny-pool
settings and query into one Scala helper that pins fair_unified.
The copy estimate summed get_slice_memory_size(), which counts the whole child array of a
sliced list. Batches sliced from one list array were each charged for all of its values, so a
window partition that fit could be refused. Reserve the copy right after concatenating instead,
counting only the buffers it does not share with the batches it was made from. That is exact
for nested, dictionary and view types.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: WindowAggExec buffered the entire task input without reserving memory.
  • Design approach: CometWindowAggExec emits completed window partitions and accounts for retained buffers and concatenation copies.
  • Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Checked partition processing against Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The existing shared-list overcounting concern is fixed at this head.
  • Key design decisions: Reusing DataFusion’s plan metadata and window evaluators limits duplication. RecordBatchMemoryCounter accounts for shared buffers without per-slice estimates.
  • Implementation sketch: Reviewed the complete five-file, three-commit diff, including the operator, planner integration, Rust and Scala tests, and tuning documentation.
  • Behavioral changes worth calling out: Completed partitions return earlier. Refused reservations fail the task because spilling remains unsupported. Boundary checks and memory accounting add per-batch work. BoundedWindowAggExec remains unchanged.
  • Suggested improvements: No additional P1/P2 changes requested.

Reviewed full SHA bde4aa0f760e1c6107a2863076c08931d28db850 against 605051ad239ef704f5f25d67910a446a6b6d7c70. The PR is non-draft. Read the snapshot and live reviews, comments and threads. The existing thread remains administratively unresolved, but its reported failure is no longer reproducible.

Routed skills: review-comet-pr, review-comet-memory-pr and review-comet-expression-pr, with the operator and related contributor guides.

Exact-head CI: 22 checks passed, 13 were skipped and the Spark 4.1 execution-suite job remained running. Rust tests, native build, other Comet Spark 4.1 suites, and TPC-H/TPC-DS verification passed. Spark SQL jobs were skipped. No failures were reported.

Validation: All six added Rust tests passed in an isolated harness importing the unchanged exact-head operator with DataFusion 55.1.0 and Arrow 59.3.0. The original 10,000-row shared-list reproduction also passed under its 16 MiB limit and matched WindowAggExec. Full local Comet/JVM and Spark SQL suites were not run because the shared disk had less than 1 GB free. Author-reported performance measurements were not independently reproduced.

@andygrove
andygrove requested a review from comphead September 27, 2026 20:52

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: WindowAggExec buffered the entire task input without reserving memory.
  • Design approach: CometWindowAggExec emits completed window partitions and reserves retained buffers and concatenation copies.
  • Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Compared partition processing and affected window semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The existing nested-list overcounting concern is fixed and its reproduction now passes.
  • Key design decisions: Reusing DataFusion’s schema, ordering, plan properties and expression evaluators limits duplication. RecordBatchMemoryCounter handles shared buffers without per-slice estimates.
  • Implementation sketch: Reviewed the complete five-file, three-commit base-relative diff, covering the operator, planner integration, Rust and Scala tests, and tuning documentation.
  • Behavioral changes worth calling out: Completed partitions return earlier. Refused reservations now fail the task because spilling remains unsupported. Partition detection, boundary-key concatenation and memory accounting add per-batch work. BoundedWindowAggExec remains unchanged.
  • Suggested improvements: No additional P1/P2 changes requested.

Reviewed full SHA bde4aa0f760e1c6107a2863076c08931d28db850 against base 605051ad239ef704f5f25d67910a446a6b6d7c70. The PR remains non-draft. Read the snapshot and live reviews, issue comments, inline comments and threads. The existing thread remains administratively unresolved, but its reported failure is no longer reproducible.

Routed skills: review-comet-pr, review-comet-memory-pr and review-comet-expression-pr, with the relevant contributor guides.

Exact-head CI: 24 checks passed and 14 were skipped, with no failures or pending checks. Rust tests, native build, Comet Spark 4.1 suites and TPC-H/TPC-DS verification passed. The execution-suite artifact confirms 1,100 passing tests, including all three new window tests. Broader Spark SQL and macOS jobs were skipped. Raw job-log access was blocked by network policy, but check results and the execution test report were inspected.

Validation: A freshly compiled standalone harness importing the unchanged exact-head operator passed all six added Rust tests and the original shared-list reproduction under its 16 MiB limit, matching WindowAggExec. Dependencies were DataFusion 55.1.0 and Arrow 59.3.0. Full Comet/JVM and Spark SQL suites were not rerun locally. Author-reported performance measurements were not independently reproduced.

@comphead

Copy link
Copy Markdown
Contributor

Thanks for this change. Emitting each window partition as soon as the sorted input moves past it matches Spark's partition-at-a-time WindowExec (minus spilling), and it is similar to how StarRocks' analytic operator processes clustered input. The differential test against WindowAggExec and the mutation-checked memory tests are convincing. A few things I'd like to see addressed.

Memory pool interplay

  1. Under the default fair_unified pool, each consumer is limited to the pool divided by the number of registered consumers (native/core/src/execution/memory_pools/fair_pool.rs:165-171). A sort then window task has three (ExternalSorter, ExternalSorterMerge and the window). The window reserves a partition and its concatenated copy together, so by my reading the largest partition that fits is about a sixth of the pool. That share is the binding limit when a task has more consumers than there are running tasks, such as a skewed straggler, which is exactly the case this PR is for. In that case, registering in execute() (native/core/src/execution/operators/window_agg.rs:191) also cuts the sort's share from pool/2 to pool/3 while it reads its input, even when every window partition is tiny. I have not measured this. A few suggestions:
    • Register the consumer when the first batch is buffered. The sort has consumed its input by then. Make fair_unified account for spillable consumers #5465 would be the real fix in the pool.
    • Have the error hint (window_agg.rs:428-431) and the tuning doc say that a partition needs roughly twice its size within its share.
    • Suggest spark.comet.exec.memoryPool=greedy_unified before disabling native windows, as docs/source/user-guide/latest/migration-guide.md:88-90 already does. spark.comet.exec.window.enabled=false also moves bounded windows back to Spark, which may be worth saying.

Code

  1. evaluate counts its inputs a second time (window_agg.rs:332-344). self.buffered_memory has already counted every buffer in batches, and concat only reuses input buffers, so self.buffered_memory.count_batch(&batch) should return the same number (and 0 for the zero-copy single-batch case). That would remove the second counter, the extra pass over every buffered batch, and the batches.len() > 1 branch with its comment.
  2. Key evaluation, evaluate_partition_ranges, the continuation check and the reservation in push_batch run outside elapsed_compute (window_agg.rs:329). Batches that only extend the open partition record no time, while WindowAggExec times its range pass. Could we time push_batch and the final evaluate in poll_next_inner and drop the inner timer, since nested timers on one Time count twice? Also, continues_buffered_partition slices every column of the last buffered batch and re-evaluates its keys on every batch (window_agg.rs:305). Keeping the last row's key arrays from line 263 would avoid that and still use the same partition kernel.

Tests

  1. A few suggestions:
    • The first new Scala test (CometWindowExecSuite.scala:1517) looks like a good fit for a SQL file test under spark/src/test/resources/sql-tests/windows/, with -- Config: spark.comet.batchSize=7 and INSERT ... SELECT ... FROM range(2000). Its second query (no PARTITION BY) could go. That path does not depend on the function, and it is already covered by "aggregate window function for all types" (OVER() at batch size 128) and by the Rust fuzz with partitioned = false.
    • Tests 2 and 3 need to stay in Scala, since expect_error requires Spark to fail as well. Test 3 (:1596) asserts DataFusion's TrackConsumersPool wording. Asserting on the new CometWindowAggExec cannot spill hint instead would also prove the hint reaches Spark. The Context text does survive into CometNativeException.
    • In the Rust tests, the interleaved empty-batch block in empty_input_and_empty_batches (window_agg.rs:727-737) repeats the fuzz, because random_slices already produces empty batches. The zero-row cases at 724-725 are the part the fuzz cannot produce. The peak assertions (750, 835) repeat what the neighboring refusals prove, so run_with_limit could return just the batches and keep its reserved == 0 check.

Tracking and docs

  1. The removal condition at window_agg.rs:38 has no open upstream work behind it. Support spilling for WindowAggExec datafusion#22947 and feat: memory accounting for WindowAggExec datafusion#23207 (reservation only, still buffering the whole input) were both closed unmerged. Would you consider offering this stream upstream as the no-spill first step of Support spilling for WindowAggExec datafusion#22946 and linking it there? Otherwise Comet carries its own copy of WindowAggStream and about ten delegated methods through every DataFusion upgrade.
  2. Native window operators reserve no memory for the batches they buffer #6253 covers both window operators, and BoundedWindowAggExec still reserves nothing. Its RANGE frames keep a whole group of peer rows, so a skewed ORDER BY value can still grow untracked. Could this be Part of #6253, or could we file a follow-up for BoundedWindowAggExec?
  3. A few docs now leave the window out:
    • docs/source/user-guide/latest/tuning/memory.md:44-46 lists what the pool tracks.
    • docs/source/contributor-guide/memory_management.md:556-558 names only ShuffledHashJoin as an operator that cannot spill.
    • docs/source/contributor-guide/memory_management.md:366-369 names only sort and hash join as reserving imported batches. OVER () has no sort below it, so the window can buffer scan or shuffle batches directly.

CI and benchmarks

  • The Spark SQL suites were skipped in PR CI. Since this changes the planner and adds a native operator, could we apply run-spark-4.1-tests before it goes to the merge queue?
  • If the standalone bench ran on a DataFusion pool, two costs are not in its numbers: extra sort spills from the share dilution above, and three JNI-backed pool calls for every batch that completes a partition. A TPC-DS window query under fair_unified, with a pool small enough to make the sort spill and compared against main, would settle both.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: WindowAggExec buffered the entire task input without reserving memory.
  • Design approach: CometWindowAggExec emits completed window partitions and reserves retained buffers and concatenation copies.
  • Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Compared partition processing against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The existing nested-list overcounting concern is fixed.
  • Key design decisions: Reusing DataFusion’s schema, ordering, plan properties and window evaluators limits duplication. RecordBatchMemoryCounter accounts for shared buffers.
  • Implementation sketch: Reviewed the complete five-file, three-commit base-relative diff, covering the operator, planner integration, Rust and Scala tests, and tuning documentation.
  • Behavioral changes worth calling out: Completed partitions return earlier. Refused reservations now fail the task because spilling remains unsupported. Partition detection and accounting add per-batch work. BoundedWindowAggExec remains unchanged.
  • Suggested improvements: No additional P1/P2 changes requested.

Reviewed full SHA bde4aa0f760e1c6107a2863076c08931d28db850 against base 605051ad239ef704f5f25d67910a446a6b6d7c70. Confirmed the PR is non-draft. Read snapshot and live reviews, issue comments, inline comments and threads. The existing thread remains administratively unresolved, but its reported failure no longer reproduces.

Routed skills: review-comet-pr, review-comet-memory-pr and review-comet-expression-pr, with relevant contributor guidance.

Exact-head CI: 24 checks passed and 14 were skipped, with no failures or pending checks. Rust tests, native build, Comet Spark 4.1 suites and TPC-H/TPC-DS verification passed. Verified the execution-report artifact’s digest against GitHub metadata. It records 1,100 passing tests, including all three added window tests.

Validation: A freshly compiled standalone harness importing the unchanged exact-head operator passed all six added Rust tests and the original shared-list reproduction under its 16 MiB limit, matching WindowAggExec. Dependencies match the lockfile: DataFusion 55.1.0 and Arrow 59.3.0. Full Comet/JVM and Spark SQL suites were not rerun locally. Broader Spark SQL and macOS CI jobs were skipped. Author-reported performance measurements were not independently reproduced.

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the thorough tests and the benchmark numbers. The differential test against WindowAggExec and the mutation checks are reassuring. Beyond the points in the earlier review comment, my only addition is the upgrade guide question inline on the new tuning section.

aggregate over `PARTITION BY` with no `ORDER BY` such as `sum(x) OVER (PARTITION BY k)`, need a whole window partition
before they return anything. Comet buffers one window partition at a time for them and reserves it from its native
memory pool, but it does not yet provide spill-to-disk for them. A window partition that does not fit in its share of
the pool, for example because of a heavily skewed key or a window without `PARTITION BY`, fails the task with a memory

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This section describes a change in outcome. A window partition that does not fit its share now fails the task, where before it ran on untracked memory. Would it make sense to add an entry for it to the upgrade guide? docs/source/contributor-guide/config_conventions.md (Changing the Behavior of an Existing Config) asks for one, plus a spark.comet.legacy.* key. I'm not sure a key is expected for a memory accounting fix, so a note like the 1.1.0 "Memory Pool Limits" entry may be enough.

@andygrove andygrove added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 30, 2026
@andygrove

Copy link
Copy Markdown
Member Author

I would like to do a deeper review of this PR before it is merged.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Native window operators reserve no memory for the batches they buffer

3 participants