andygrove opened a new pull request, #6279:
URL: https://github.com/apache/datafusion-comet/pull/6279

   ## 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 `WindowExpr`s 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 before it is made, as 
apache/datafusion#25496 does for the hash join build side. 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, so its memory is bounded by the frame rather than the window 
partition. 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.** Unchanged for 10-row window partitions, and 3–18% slower 
elsewhere. That is at most 0.6 ms per 2M rows, spent on a second 
partition-boundary pass per batch.
   - **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.
     - Each of four mutations fails at least one of these tests: no batch 
reservation, no copy reservation, no release after emitting, and never 
continuing a window partition across batches.
   - **`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`.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to