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]