SubhamSinghal opened a new pull request, #25968:
URL: https://github.com/apache/datafusion/pull/25968

   ## Which issue does this PR close?
   
   Related to #6899 (WindowTopN). Follows #25730 (merged), which applied this 
design to `ROW_NUMBER`. This PR brings `RANK` to the same design.
   
   ## Rationale for this change
   
   With `enable_window_topn`, `RANK() ... WHERE rk <= K` on high-cardinality 
`PARTITION BY` is **slower than the sort plan it replaces**: at 100K partitions 
the rewrite costs 2.76× the sort plan's CPU. 
`PartitionedTopKRank::insert_batch`:
   
   - gathers one `take_record_batch` sub-batch per (input batch, partition) 
pair, roughly one small `RecordBatch` per input row at 100K partitions;
   - does a 1-row `take_record_batch` for every row evicted into the tie list;
   - runs a full `TopKHeap`, with its own `RecordBatchStore`, per partition;
   - sums `size()` over every partition on every batch.
   
   Profiling main at 1M partitions put about 58% of samples in creating and 
destroying those small batches. Payload width made almost no difference, so the 
cost is the number of calls, not the bytes copied. #25730 removed exactly this 
for `ROW_NUMBER`. This PR does the same for `RANK`, plus boundary ties.
   
   ## What changes are included in this PR?
   
   In `datafusion/physical-plan/src/topk/mod.rs`, plus a doc update in 
`sorts/partitioned_topk.rs`:
   
   - **`PartitionedTopKRank` uses #25730's design.** It encodes the partition 
and ORDER BY columns once per batch, keeps or drops each row in its partition's 
`PartitionHeap` in one pass, then gathers **the admitted rows once per input 
batch** into an operator-wide `RecordBatchStore`. Output is interleaved out of 
that store in `batch_size` chunks by `EmitState`, so the `BatchCoalescer` goes 
away.
   - **Ties are 8-byte `StoreRef`s** (`batch_id`, `row`) in a per-partition 
`ties` list. Every tie shares the heap root's key, so none is stored.
     - Evicted but still tied: the row's reference and its store use move from 
the heap to the tie list. No 1-row `take_record_batch`.
     - Cutoff improves: the evicted row and all ties are released, one `unuse` 
each, so amortised O(1) per admitted row.
   - **Memory stays proportional to rows retained** (the property #24591 
established). The store holds only admitted rows, and the same ratio-2 
compaction as #25730 keeps `store.total_rows ≤ 2 × live_slots` after every 
batch, where `live_slots` counts heap rows and ties. Ties from every partition 
a batch touches share that batch's one store entry, so they are charged once, 
not once per partition (the over-count in #23326).
   - **`size()` is O(1)**, from running totals plus `store.size()`.
   - **Shared with `ROW_NUMBER`.** Compaction (`RecordBatchStore::compact`, 
replacing `PartitionedTopK::compact_store`), the in-flight batch's bookkeeping 
(`pending` / `release` / `insert_rows`) and stream construction 
(`EmitState::stream`) are now used by both operators. No behaviour change for 
`ROW_NUMBER`.
   - **`EvictedRow` is removed.** `RANK` was its last reader. `TopKHeap::add` 
no longer looks up and clones the evicted row's batch.
   
   ## Are these changes tested?
   
   Yes. All existing `PartitionedTopKRank` tests pass, including the three 
memory-bound tests from #24591. New tests:
   
   - `test_partitioned_topk_rank_bookkeeping_tracks_recompute`: 256 randomized 
shapes, checking after every batch that each store entry's `uses` matches the 
rows pointing into it, `total_rows ≤ 2 × live_slots`, the running totals and 
reservation match a recompute, and the output matches brute force in `(pk, 
val)` order. It replaces `test_partitioned_topk_rank_matches_bruteforce`.
   - `test_partitioned_topk_rank_store_bounded_when_ties_spread_thinly`, 
`test_partitioned_topk_rank_boundary_move_releases_ties_across_batches`, and 
`test_partitioned_topk_rank_ties_share_one_store_entry` (the #23326 shape).
   
   ## Benchmarks
   
   `benchmarks/queries/h2o/window.sql` RANK queries Q18–Q23 (h2o J1 `large`, 
10M rows, K=2). User CPU, median of 7 alternating runs, 14 cores. **base** = 
#25730's tip, whose `RANK` code is main's; **sort plan** = flag off.
   
   | partitions | main, ON vs OFF | this branch, ON vs OFF |
   |---|---|---|
   | 100 (Q18) | 5.87× faster | 6.07× faster |
   | 1,000 (Q19) | 5.31× faster | 6.00× faster |
   | 1,000 (Q20) | 4.78× faster | 5.24× faster |
   | 10,000 (Q21) | 2.02× faster | 5.44× faster |
   | 10,000 (Q22) | 2.19× faster | 4.74× faster |
   | 100,000 (Q23) | 2.76× slower | 1.83× faster |
   
   ## Are there any user-facing changes?
   
   No. `enable_window_topn` stays default-false,


-- 
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