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]

Reply via email to