Conversation
…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.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
WindowAggExecretained the entire task input without memory reservations. - Design approach:
CometWindowAggExecemits 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.
BoundedWindowAggExecremains 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()?; |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
WindowAggExecbuffered the entire task input without reserving memory. - Design approach:
CometWindowAggExecemits 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.
RecordBatchMemoryCounteraccounts 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.
BoundedWindowAggExecremains 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.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
WindowAggExecbuffered the entire task input without reserving memory. - Design approach:
CometWindowAggExecemits 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.
RecordBatchMemoryCounterhandles 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.
BoundedWindowAggExecremains 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.
|
Thanks for this change. Emitting each window partition as soon as the sorted input moves past it matches Spark's partition-at-a-time Memory pool interplay
Code
Tests
Tracking and docs
CI and benchmarks
|
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
WindowAggExecbuffered the entire task input without reserving memory. - Design approach:
CometWindowAggExecemits 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.
RecordBatchMemoryCounteraccounts 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.
BoundedWindowAggExecremains 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
left a comment
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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.
|
I would like to do a deeper review of this PR before it is merged. |
Which issue does this PR close?
Closes #6253.
Rationale for this change
PERCENT_RANK,CUME_DIST,NTILE, and aggregates whose frame ends atUNBOUNDED FOLLOWINGrun in DataFusion'sWindowAggExec. That includes Spark's default frame foragg(x) OVER (PARTITION BY k).WindowAggExecbuffers 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.
WindowAggExecholds 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 failedSUM(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?
CometWindowAggExecwhere it usedWindowAggExec. It evaluates the sameWindowExprs one window partition at a time, asWindowAggExecdoes, with two differences:PARTITION BYkeys, 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. WithoutPARTITION BY, it still buffers the whole input.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.BoundedWindowAggExecis unchanged. It keeps only the rows its frames still need, which is bounded by the frame rather than the window partition (for aRANGE ... CURRENT ROWframe, that is a whole group of peer rows). It still reserves nothing.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
WindowAggExecis 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
WindowAggExecwith a standalone bench over 2M rows in 8192-row batches: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.WindowAggExecCometWindowAggExecWithout
PARTITION BY, both peak at 1,090 MB, butCometWindowAggExecreserves 960 MB of it.How are these changes tested?
Rust unit tests in
window_agg.rs.WindowAggExecover 20 seeds. It uses random window partition sizes, random batch boundaries including empty batches, null keys and two key columns, with and withoutPARTITION BY.get_slice_memory_sizeestimate over-charged this case about tenfold and was refused.CometWindowExecSuite. The new tests:Additional allocation failed for CometWindowAggExec.The suite passes on Spark 3.4, 3.5 and 4.1.
Other suites.
DataFrameWindowFunctionsSuiteandDataFrameWindowFramesSuiteran 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 marksIgnoreCometor adjusts forCometWindowExec.