SubhamSinghal commented on code in PR #25968:
URL: https://github.com/apache/datafusion/pull/25968#discussion_r4182848096
##########
datafusion/physical-plan/src/topk/mod.rs:
##########
@@ -1313,6 +1314,145 @@ impl RecordBatchStore {
+ self.batches.capacity() * (size_of::<u32>() +
size_of::<RecordBatchEntry>())
+ self.batches_size
}
+
+ /// Rewrite the store to hold only the rows `partitions` still reference,
+ /// once it holds [`STORE_COMPACTION_RATIO`]× more than the `live_slots` of
+ /// them, and repoint every partition at the rows' new places.
+ /// [`PartitionedTopK`] and [`PartitionedTopKRank`] call this at the end of
+ /// every `insert_batch`.
+ ///
+ /// An entry holds every row *admitted* from its input batch and is freed
+ /// only when the last of them is evicted, so when survivors spread thinly
+ /// one live row keeps a whole entry resident and residency tracks the
+ /// *input*, not `partitions × K`: 512 partitions of `K = 1` fed 512
batches
+ /// pin 131 K rows to retain 512. Nothing else bounds that. A single entry
is
+ /// no exception — rows admitted then superseded within their own batch
stay
+ /// in the gather unreferenced — so this does not skip a one-entry store.
+ /// `RANK` ties add a second way to leave rows unreferenced, since a
+ /// boundary move releases every tie of the partition at once.
+ ///
+ /// Amortized O(1) per admitted row: one pass over the live slots, and it
+ /// cannot recur until the store has taken on another `live_slots` rows.
+ ///
+ /// Rows are rewritten into `batch_size` chunks, not one batch: an entry is
+ /// released only when its last slot is evicted, so a single batch of every
+ /// live row would free nothing until every partition has churned.
+ ///
+ /// Peak residency is the old store plus the new one — chunks are built
+ /// before the old entries drop, and the reservation is not resized until
+ /// `insert_batch` returns — so a pool sized at the steady-state bound can
be
+ /// exceeded transiently without erroring.
+ ///
+ /// All-or-nothing: plan the move and interleave, which is the only
fallible
+ /// step, before rewriting the slots and the store. A failing interleave
+ /// leaves the operator as it was rather than holding slots pointing at ids
+ /// the store never got.
+ fn compact<P: StoreSlots>(
+ &mut self,
+ partitions: &mut HashMap<Vec<u8>, P>,
+ live_slots: usize,
+ batch_size: usize,
+ ) -> Result<()> {
+ if self.total_rows <= live_slots * STORE_COMPACTION_RATIO {
+ return Ok(());
+ }
+
+ // Scoped so these clones, which keep the old batches alive while the
+ // 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();
+
+ // Keyed by the row's current place rather than by its position in
+ // this walk, so the partitions are free to repoint in any order —
+ // 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> =
Review Comment:
Addressed in f9d7919e1fed26f901ab6b5a2730e0057b85d56d
--
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]