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]