perf: PartitionedTopKRank on the shared store with decide-then-gather - #25968
SubhamSinghal wants to merge 4 commits into
Conversation
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #25968 +/- ##
==========================================
+ Coverage 82.60% 82.67% +0.07%
==========================================
Files 1145 1147 +2
Lines 444057 446767 +2710
Branches 444057 446767 +2710
==========================================
+ Hits 366796 369348 +2552
+ Misses 54996 54995 -1
- Partials 22265 22424 +159 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
@kosiew @jayzhan211 @kumarUjjawal can you help in reviewing this PR? |
kumarUjjawal
left a comment
There was a problem hiding this comment.
Thank you @SubhamSinghal
Left few comments please take a look.
| released_slots += state.ties.len(); | ||
| state.ties.clear(); |
There was a problem hiding this comment.
Release tie-list capacity when the boundary improves
| EmitState::stream( | ||
| schema, | ||
| futures::stream::iter(out), | ||
| ))) | ||
| metrics, | ||
| reservation, | ||
| batch_size, | ||
| &store, |
There was a problem hiding this comment.
Limit interleave inputs to batches referenced by each output chunk
| match state.heap.classify(k, key) { | ||
| // Strictly worse than the boundary: drop the row. | ||
| Some(Ordering::Greater) => continue, | ||
| Some(Ordering::Equal) => { |
There was a problem hiding this comment.
In PartitionedTopKRank::insert_batch, rows admitted as ties (the Equal arm) do not increment replacements, but every other admission does, and every admission does in PartitionedTopK, so row_replacements undercounts rows the operator retained.
| // no two slots share a store row, so the key identifies exactly one | ||
| // slot. | ||
| let mut coords: Vec<(usize, usize)> = Vec::with_capacity(live_slots); | ||
| let mut moved: HashMap<StoreRef, StoreRef> = |
There was a problem hiding this comment.
RecordBatchStore::compact builds a HashMap<StoreRef, StoreRef> with live_slots entries just to repoint slots, although partitions.values()/values_mut() and store_rows()/repoint() walk the slots in the same order, so a running counter would give each slot its new place.
| // compacted ones are built, drop before the store is rewritten. | ||
| let first_id = self.next_batch_id(); | ||
| let (moved, chunks) = { | ||
| let (old, array_pos) = self.positional(); |
There was a problem hiding this comment.
compact calls positional(), which clones every stored RecordBatch (an Arc bump per column) and builds a fresh HashMap, but compact only needs borrowed &RecordBatch refs that it drops before the store is rewritten.
| Some(Ordering::Less) => { | ||
| // Replacing the root overwrites its key in place, so | ||
| // keep a copy to tell whether the boundary moved. | ||
| evicted_key.clear(); |
There was a problem hiding this comment.
Every strictly-better eviction in RANK memcpys the old root key into evicted_key only to compare it with the new root. PartitionHeap::add could hand the old key back by swapping Vecs instead.
|
|
||
| // 1. Evaluate + encode partition columns into the reusable | ||
| // scratch (cleared then appended). | ||
| // 1. Evaluate the partition and ORDER BY columns and encode each once |
There was a problem hiding this comment.
PartitionedTopKRank still duplicates PartitionedTopK almost line for line: the same ~14 fields, try_new, the phase-1 evaluate/encode block, the phase-3 insert_rows/compact/try_resize tail, the emit destructuring and size(). Only the per-row admission rule differs.
There was a problem hiding this comment.
PR adds RecordBatchStore::positional to build the batches/batch_id -> position pair, but TopKHeap::emit_with_state still builds the same mapping inline.
There was a problem hiding this comment.
Comments left stale by the rewrite: the doc and step comments of test_partitioned_topk_rank_boundary_shifts_clears_ties still describe the removed in-flight equal_indices list ('must clear both state.ties and the in-flight equal_indices'), and assert_store_matches_slots (line 4948) still names compact_store, which this PR replaced with RecordBatchStore::compact.
…artitioned TopK operators
jayzhan211
left a comment
There was a problem hiding this comment.
Thanks @SubhamSinghal! Nice work — good to go from my side.
Which issue does this PR close?
Part of #6899. Follows #25730 (merged), which applied this design to
ROW_NUMBER. This PR bringsRANKto the same design.Rationale for this change
With
enable_window_topn,RANK() ... WHERE rk <= Kon high-cardinalityPARTITION BYis slower than the sort plan it replaces: at 100K partitions the rewrite costs 2.76× the sort plan's CPU.PartitionedTopKRank::insert_batch:take_record_batchsub-batch per (input batch, partition) pair, roughly one smallRecordBatchper input row at 100K partitions;take_record_batchfor every row evicted into the tie list;TopKHeap, with its ownRecordBatchStore, per partition;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 forRANK, plus boundary ties.What changes are included in this PR?
In
datafusion/physical-plan/src/topk/mod.rs, plus a doc update insorts/partitioned_topk.rs:PartitionedTopKRankuses perf: PartitionedTopK hold retained rows in one shared RecordBatchStore #25730's design. It encodes the partition and ORDER BY columns once per batch, keeps or drops each row in its partition'sPartitionHeapin one pass, then gathers the admitted rows once per input batch into an operator-wideRecordBatchStore. Output is interleaved out of that store inbatch_sizechunks byEmitState, so theBatchCoalescergoes away.StoreRefs (batch_id,row) in a per-partitiontieslist. Every tie shares the heap root's key, so none is stored.take_record_batch.unuseeach, so amortised O(1) per admitted row.store.total_rows ≤ 2 × live_slotsafter every batch, wherelive_slotscounts 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 PartitionedTopKRank over-accounts memory when a single batch has boundary ties across many partitions #23326).size()is O(1), from running totals plusstore.size().ROW_NUMBER. Compaction (RecordBatchStore::compact, replacingPartitionedTopK::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 forROW_NUMBER.EvictedRowis removed.RANKwas its last reader.TopKHeap::addno longer looks up and clones the evicted row's batch.Are these changes tested?
Yes. All existing
PartitionedTopKRanktests 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'susesmatches 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 replacestest_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, andtest_partitioned_topk_rank_ties_share_one_store_entry(the PartitionedTopKRank over-accounts memory when a single batch has boundary ties across many partitions #23326 shape).Benchmarks
benchmarks/queries/h2o/window.sqlRANK queries Q18–Q23 (h2o J1large, 10M rows, K=2). User CPU, median of 7 alternating runs, 14 cores. base = #25730's tip, whoseRANKcode is main's; sort plan = flag off.Are there any user-facing changes?
No.
enable_window_topnstays default-false,