comphead commented on PR #6279:
URL:
https://github.com/apache/datafusion-comet/pull/6279#issuecomment-5874119437
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. #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**
2. `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.
3. 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**
4. 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**
5. The removal condition at `window_agg.rs:38` has no open upstream work
behind it. apache/datafusion#22947 and apache/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
apache/datafusion#22946 and linking it there? Otherwise Comet carries its own
copy of `WindowAggStream` and about ten delegated methods through every
DataFusion upgrade.
6. #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`?
7. 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.
--
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]