kumarUjjawal commented on code in PR #23869:
URL: https://github.com/apache/datafusion/pull/23869#discussion_r3672678798
##########
datafusion/physical-plan/src/topk/mod.rs:
##########
@@ -1811,6 +1811,376 @@ impl PartitionedTopKRank {
}
}
+/// Per-partition state for `DENSE_RANK()` semantics.
+///
+/// A `HashMap<Vec<u8>, Vec<TieEntry>>` keyed by the row-encoded ORDER BY
+/// bytes, capped at `k` distinct keys. Each key's `Vec<TieEntry>` holds
+/// every row seen at that ob value, one entry per contributing source
+/// `RecordBatch`. `TieEntry` is reused verbatim from [`RankPartitionState`].
+struct DenseRankPartitionState {
+ groups: HashMap<Vec<u8>, Vec<TieEntry>>,
+}
+
+impl DenseRankPartitionState {
+ fn size(&self) -> usize {
+ let table_overhead =
+ self.groups.capacity() * (size_of::<Vec<u8>>() +
size_of::<Vec<TieEntry>>());
+ let contents: usize = self
+ .groups
+ .iter()
+ .map(|(key, entries)| {
+ key.capacity()
+ + entries.capacity() * size_of::<TieEntry>()
+ + entries
+ .iter()
+ .map(|e| {
+ e.row_indices.capacity() * size_of::<u32>() +
e.batch_bytes
+ })
+ .sum::<usize>()
+ })
+ .sum();
+ table_overhead + contents
+ }
+}
+
+/// Sibling to [`PartitionedTopK`] / [`PartitionedTopKRank`] implementing
+/// `DENSE_RANK()` semantics.
+///
+/// Per partition, retains every row whose ORDER BY value is among the K
+/// distinct-smallest ob values seen for that partition. The total row
+/// count kept per partition is unbounded in `rows_per_distinct_value`
+/// (unlike `RANK`, which is bounded above by K + boundary ties).
+///
+/// Like [`PartitionedTopK`], the [`RowConverter`], [`MemoryReservation`],
+/// scratch [`Rows`] buffer, and [`TopKMetrics`] are shared across all
+/// partitions for this operator instance.
+///
+/// # Algorithm (per batch)
+///
+/// Evaluate + encode partition-by and order-by columns once, then group
+/// the batch's row indices by partition key. For each partition, bucket
+/// that partition's rows by distinct ob value (a within-call
+/// accumulation), then merge each bucket into the partition state. Every
+/// bucket is built from the current batch's rows, so each `TieEntry` is
+/// pinned to the batch its `row_indices` point into.
+///
+/// For each partition, for each distinct `ob_key` run in this batch:
+/// - `ob_key` already in `state.groups` → push this batch's run as a
+/// new `TieEntry` (one entry per contributing batch).
+/// - `ob_key` new, `state.groups.len() < k` → insert the run as a new
+/// group.
+/// - `ob_key` new, `state.groups.len() == k` → find the current max via
+/// an O(K) scan of `state.groups.keys()`:
+/// - `ob_key < max` → remove the max key (evict the entire max-key
+/// group — up to many rows) and insert the run. The evicted group's
+/// row count is added to the `row_replacements` metric.
+/// - `ob_key >= max` → drop the whole run; no map mutation.
+pub(crate) struct PartitionedTopKDenseRank {
+ schema: SchemaRef,
+ metrics: TopKMetrics,
+ reservation: MemoryReservation,
+ /// ORDER BY expressions (excludes PARTITION BY).
+ expr: LexOrdering,
+ /// Encoder for ORDER BY columns. Reused across partitions.
+ row_converter: RowConverter,
+ /// Scratch row buffer reused across `insert_batch` calls.
+ scratch_rows: Rows,
+ /// PARTITION BY expressions.
+ partition_exprs: Vec<Arc<dyn PhysicalExpr>>,
+ /// Encoder for the partition key.
+ partition_converter: RowConverter,
+ /// Scratch row buffer for partition-key encoding. Reused across
+ /// `insert_batch` calls (cleared + appended each batch).
+ partition_scratch_rows: Rows,
+ /// One state per distinct partition key seen so far. Keyed by the
+ /// row-encoded PARTITION BY bytes (byte-comparable encoding, so the
+ /// `Vec<u8>` hashes, compares, and sorts identically to an
+ /// `OwnedRow`) which lets `insert_batch` look partitions up with
+ /// `entry_ref` — allocating a key only on first sight of a partition
+ /// rather than once per row.
+ states: HashMap<Vec<u8>, DenseRankPartitionState>,
+ /// Scratch map reused across `insert_batch` calls to group a batch's
+ /// row indices by partition key. Drained (not reallocated) each batch
+ /// so its backing table is allocated once, not per batch.
+ partition_groups: HashMap<Vec<u8>, Vec<u32>>,
+ /// Scratch map reused across partitions within a batch to bucket a
+ /// partition's rows by distinct ORDER BY value. Drained (not
+ /// reallocated) per partition so its backing table is allocated once,
+ /// not once per distinct partition key.
+ ob_runs: HashMap<Vec<u8>, Vec<u32>>,
+ k: usize,
+ batch_size: usize,
+}
+
+impl PartitionedTopKDenseRank {
+ #[expect(clippy::too_many_arguments)]
+ pub(crate) fn try_new(
+ partition_id: usize,
+ schema: SchemaRef,
+ partition_exprs: Vec<Arc<dyn PhysicalExpr>>,
+ partition_sort_fields: Vec<SortField>,
+ order_expr: LexOrdering,
+ k: usize,
+ batch_size: usize,
+ runtime: &Arc<RuntimeEnv>,
+ metrics: &ExecutionPlanMetricsSet,
+ ) -> Result<Self> {
+ assert!(k > 0, "PartitionedTopKDenseRank requires k > 0");
+ let reservation =
+
MemoryConsumer::new(format!("PartitionedTopKDenseRank[{partition_id}]"))
+ .register(&runtime.memory_pool);
+
+ let order_sort_fields = build_sort_fields(&order_expr, &schema)?;
+ let row_converter = RowConverter::new(order_sort_fields)?;
+ let scratch_rows =
+ row_converter.empty_rows(batch_size, ESTIMATED_BYTES_PER_ROW *
batch_size);
+
+ let partition_converter = RowConverter::new(partition_sort_fields)?;
+ let partition_scratch_rows = partition_converter
+ .empty_rows(batch_size, ESTIMATED_BYTES_PER_ROW * batch_size);
+
+ Ok(Self {
+ schema,
+ metrics: TopKMetrics::new(metrics, partition_id),
+ reservation,
+ expr: order_expr,
+ row_converter,
+ scratch_rows,
+ partition_exprs,
+ partition_converter,
+ partition_scratch_rows,
+ states: HashMap::new(),
+ partition_groups: HashMap::new(),
+ ob_runs: HashMap::new(),
+ k,
+ batch_size,
+ })
+ }
+
+ /// Encode PARTITION BY and ORDER BY columns once, demultiplex the
+ /// batch's rows by partition key, then per partition bucket the rows
+ /// by distinct ob value and merge each bucket into the partition
+ /// state as one [`TieEntry`].
+ pub(crate) fn insert_batch(&mut self, batch: &RecordBatch) -> Result<()> {
+ let baseline = self.metrics.baseline.clone();
+ let _timer = baseline.elapsed_compute().timer();
+
+ let num_rows = batch.num_rows();
+ if num_rows == 0 {
+ return Ok(());
+ }
+
+ // Captured once so every `TieEntry` push from this batch can
+ // reuse it (avoids `get_record_batch_memory_size` per push).
+ let input_batch_bytes = get_record_batch_memory_size(batch);
+
+ // 1. Encode partition columns.
+ let pk_arrays: Vec<ArrayRef> = self
+ .partition_exprs
+ .iter()
+ .map(|e| e.evaluate(batch).and_then(|v| v.into_array(num_rows)))
+ .collect::<Result<_>>()?;
+ self.partition_scratch_rows.clear();
+ self.partition_converter
+ .append(&mut self.partition_scratch_rows, &pk_arrays)?;
+
+ // 2. Group this batch's row indices by partition key.
+ // `partition_groups` is a reused scratch map: taken out here
+ // and drained below, so its backing table is allocated once
+ // for the operator, not once per batch. `entry_ref` owns the
+ // key only on Vacant, so it allocates one `Vec<u8>` per
+ // distinct partition rather than one per row.
+ let mut groups = std::mem::take(&mut self.partition_groups);
+ groups.clear();
+ {
+ let pk_rows = &self.partition_scratch_rows;
+ for i in 0..num_rows {
+ groups
+ .entry_ref(pk_rows.row(i).as_ref())
+ .or_default()
+ .push(i as u32);
+ }
+ }
+
+ // 3. Evaluate ORDER BY columns and encode ONCE.
+ let ob_arrays: Vec<ArrayRef> = self
+ .expr
+ .iter()
+ .map(|e| e.expr.evaluate(batch).and_then(|v|
v.into_array(num_rows)))
+ .collect::<Result<_>>()?;
+ self.scratch_rows.clear();
+ self.row_converter
+ .append(&mut self.scratch_rows, &ob_arrays)?;
+
+ let k = self.k;
+ let mut replacements: usize = 0;
+
+ // 4. Per-partition: bucket this batch's rows by distinct ob value
+ // (within-call accumulation), then merge each bucket into the
+ // partition state as a single `TieEntry`.
+ for (pk, indices) in groups.drain() {
+ let state =
+ self.states
+ .entry(pk)
+ .or_insert_with(|| DenseRankPartitionState {
+ groups: HashMap::new(),
+ });
+
+ // Bucket by ob key. `ob_runs` is a reused scratch map (taken
+ // out and drained below) so its backing table is allocated
+ // once, not once per distinct partition key. `entry_ref` owns
+ // the key only on Vacant, so repeated rows of the same ob
+ // value don't re-allocate.
+ let mut runs = std::mem::take(&mut self.ob_runs);
+ runs.clear();
+ for &orig_idx in &indices {
+ let ob_row = self.scratch_rows.row(orig_idx as usize);
+ runs.entry_ref(ob_row.as_ref()).or_default().push(orig_idx);
+ }
+
+ for (ob_key, run_indices) in runs.drain() {
+ // Case A: ob already tracked — push this batch's run as a
+ // new `TieEntry` (one entry per contributing batch, exactly
+ // like RANK pushing one tie entry per batch).
+ if let Some(entries) = state.groups.get_mut(&ob_key) {
+ entries.push(TieEntry {
+ batch: batch.clone(),
+ row_indices: run_indices,
+ batch_bytes: input_batch_bytes,
Review Comment:
that will leaves the normal dense-rank path counting one shared batch once
per retained group across every partition. This can greatly inflate the
reservation and cause false ResourcesExhausted errors. Could we fix the
shared-batch accounting here, or keep this rewrite disabled until #23326 is
fixed?
##########
datafusion/physical-plan/src/topk/mod.rs:
##########
@@ -1811,6 +1811,394 @@ impl PartitionedTopKRank {
}
}
+/// Per-partition state for `DENSE_RANK()` semantics.
+///
+/// A `HashMap<Vec<u8>, Vec<TieEntry>>` keyed by the row-encoded ORDER BY
+/// bytes, capped at `k` distinct keys. Each key's `Vec<TieEntry>` holds
+/// every row seen at that ob value, one entry per contributing source
+/// `RecordBatch`. `TieEntry` is reused verbatim from [`RankPartitionState`].
+struct DenseRankPartitionState {
+ groups: HashMap<Vec<u8>, Vec<TieEntry>>,
+ /// Cached admission boundary: the largest tracked ob value once
+ /// `groups` is full. A `HashMap` is unordered, so finding it otherwise
+ /// costs an O(K) scan of the keys on *every* new distinct ob value —
+ /// O(N·K) overall on mostly-distinct input. Caching makes the common
+ /// "new ob is worse than the boundary → drop" path O(1); the O(K) scan
+ /// is paid only when the cache is stale (`None`) — after a fill or an
+ /// eviction changes the key set.
+ max_key: Option<Vec<u8>>,
+}
+
+impl DenseRankPartitionState {
+ fn size(&self) -> usize {
+ let table_overhead =
+ self.groups.capacity() * (size_of::<Vec<u8>>() +
size_of::<Vec<TieEntry>>());
+ let contents: usize = self
+ .groups
+ .iter()
+ .map(|(key, entries)| {
+ key.capacity()
+ + entries.capacity() * size_of::<TieEntry>()
+ + entries
+ .iter()
+ .map(|e| {
+ e.row_indices.capacity() * size_of::<u32>() +
e.batch_bytes
+ })
+ .sum::<usize>()
+ })
+ .sum();
+ table_overhead + contents + self.max_key.as_ref().map_or(0, |k|
k.capacity())
+ }
+}
+
+/// Sibling to [`PartitionedTopK`] / [`PartitionedTopKRank`] implementing
+/// `DENSE_RANK()` semantics.
+///
+/// Per partition, retains every row whose ORDER BY value is among the K
+/// distinct-smallest ob values seen for that partition. The total row
+/// count kept per partition is unbounded in `rows_per_distinct_value`
+/// (unlike `RANK`, which is bounded above by K + boundary ties).
+///
+/// Like [`PartitionedTopK`], the [`RowConverter`], [`MemoryReservation`],
+/// scratch [`Rows`] buffer, and [`TopKMetrics`] are shared across all
+/// partitions for this operator instance.
+///
+/// # Algorithm (per batch)
+///
+/// Evaluate + encode partition-by and order-by columns once, then group
+/// the batch's row indices by partition key. For each partition, bucket
+/// that partition's rows by distinct ob value (a within-call
+/// accumulation), then merge each bucket into the partition state. Every
+/// bucket is built from the current batch's rows, so each `TieEntry` is
+/// pinned to the batch its `row_indices` point into.
+///
+/// For each partition, for each distinct `ob_key` run in this batch:
+/// - `ob_key` already in `state.groups` → push this batch's run as a
+/// new `TieEntry` (one entry per contributing batch).
+/// - `ob_key` new, `state.groups.len() < k` → insert the run as a new
+/// group.
+/// - `ob_key` new, `state.groups.len() == k` → the largest tracked ob
+/// value is the admission boundary, cached in `max_key` (recomputed via
+/// an O(K) scan of `state.groups.keys()` only when the cache is stale):
+/// - `ob_key < max` → remove the max key (evict the entire max-key
+/// group — up to many rows) and insert the run. The evicted group's
+/// row count is added to the `row_replacements` metric.
+/// - `ob_key >= max` → drop the whole run; no map mutation.
+pub(crate) struct PartitionedTopKDenseRank {
+ schema: SchemaRef,
+ metrics: TopKMetrics,
+ reservation: MemoryReservation,
+ /// ORDER BY expressions (excludes PARTITION BY).
+ expr: LexOrdering,
+ /// Encoder for ORDER BY columns. Reused across partitions.
+ row_converter: RowConverter,
+ /// Scratch row buffer reused across `insert_batch` calls.
+ scratch_rows: Rows,
+ /// PARTITION BY expressions.
+ partition_exprs: Vec<Arc<dyn PhysicalExpr>>,
+ /// Encoder for the partition key.
+ partition_converter: RowConverter,
+ /// Scratch row buffer for partition-key encoding. Reused across
+ /// `insert_batch` calls (cleared + appended each batch).
+ partition_scratch_rows: Rows,
+ /// One state per distinct partition key seen so far. Keyed by the
+ /// row-encoded PARTITION BY bytes (byte-comparable encoding, so the
+ /// `Vec<u8>` hashes, compares, and sorts identically to an
+ /// `OwnedRow`) which lets `insert_batch` look partitions up with
+ /// `entry_ref` — allocating a key only on first sight of a partition
+ /// rather than once per row.
+ states: HashMap<Vec<u8>, DenseRankPartitionState>,
+ /// Scratch map reused across `insert_batch` calls to group a batch's
+ /// row indices by partition key. Drained (not reallocated) each batch
+ /// so its backing table is allocated once, not per batch.
+ partition_groups: HashMap<Vec<u8>, Vec<u32>>,
+ /// Scratch map reused across partitions within a batch to bucket a
+ /// partition's rows by distinct ORDER BY value. Drained (not
+ /// reallocated) per partition so its backing table is allocated once,
+ /// not once per distinct partition key.
+ ob_runs: HashMap<Vec<u8>, Vec<u32>>,
+ k: usize,
+ batch_size: usize,
+}
+
+impl PartitionedTopKDenseRank {
+ #[expect(clippy::too_many_arguments)]
+ pub(crate) fn try_new(
+ partition_id: usize,
+ schema: SchemaRef,
+ partition_exprs: Vec<Arc<dyn PhysicalExpr>>,
+ partition_sort_fields: Vec<SortField>,
+ order_expr: LexOrdering,
+ k: usize,
+ batch_size: usize,
+ runtime: &Arc<RuntimeEnv>,
+ metrics: &ExecutionPlanMetricsSet,
+ ) -> Result<Self> {
+ assert!(k > 0, "PartitionedTopKDenseRank requires k > 0");
+ let reservation =
+
MemoryConsumer::new(format!("PartitionedTopKDenseRank[{partition_id}]"))
+ .register(&runtime.memory_pool);
+
+ let order_sort_fields = build_sort_fields(&order_expr, &schema)?;
+ let row_converter = RowConverter::new(order_sort_fields)?;
+ let scratch_rows =
+ row_converter.empty_rows(batch_size, ESTIMATED_BYTES_PER_ROW *
batch_size);
+
+ let partition_converter = RowConverter::new(partition_sort_fields)?;
+ let partition_scratch_rows = partition_converter
+ .empty_rows(batch_size, ESTIMATED_BYTES_PER_ROW * batch_size);
+
+ Ok(Self {
+ schema,
+ metrics: TopKMetrics::new(metrics, partition_id),
+ reservation,
+ expr: order_expr,
+ row_converter,
+ scratch_rows,
+ partition_exprs,
+ partition_converter,
+ partition_scratch_rows,
+ states: HashMap::new(),
+ partition_groups: HashMap::new(),
+ ob_runs: HashMap::new(),
+ k,
+ batch_size,
+ })
+ }
+
+ /// Encode PARTITION BY and ORDER BY columns once, demultiplex the
+ /// batch's rows by partition key, then per partition bucket the rows
+ /// by distinct ob value and merge each bucket into the partition
+ /// state as one [`TieEntry`].
+ pub(crate) fn insert_batch(&mut self, batch: &RecordBatch) -> Result<()> {
+ let baseline = self.metrics.baseline.clone();
+ let _timer = baseline.elapsed_compute().timer();
+
+ let num_rows = batch.num_rows();
+ if num_rows == 0 {
+ return Ok(());
+ }
+
+ // Captured once so every `TieEntry` push from this batch can
+ // reuse it (avoids `get_record_batch_memory_size` per push).
+ let input_batch_bytes = get_record_batch_memory_size(batch);
+
+ // 1. Encode partition columns.
+ let pk_arrays: Vec<ArrayRef> = self
+ .partition_exprs
+ .iter()
+ .map(|e| e.evaluate(batch).and_then(|v| v.into_array(num_rows)))
+ .collect::<Result<_>>()?;
+ self.partition_scratch_rows.clear();
+ self.partition_converter
+ .append(&mut self.partition_scratch_rows, &pk_arrays)?;
+
+ // 2. Group this batch's row indices by partition key.
+ // `partition_groups` is a reused scratch map: taken out here
+ // and drained below, so its backing table is allocated once
+ // for the operator, not once per batch. `entry_ref` owns the
+ // key only on Vacant, so it allocates one `Vec<u8>` per
+ // distinct partition rather than one per row.
+ let mut groups = std::mem::take(&mut self.partition_groups);
+ groups.clear();
+ {
+ let pk_rows = &self.partition_scratch_rows;
+ for i in 0..num_rows {
+ groups
+ .entry_ref(pk_rows.row(i).as_ref())
+ .or_default()
+ .push(i as u32);
+ }
+ }
+
+ // 3. Evaluate ORDER BY columns and encode ONCE.
+ let ob_arrays: Vec<ArrayRef> = self
+ .expr
+ .iter()
+ .map(|e| e.expr.evaluate(batch).and_then(|v|
v.into_array(num_rows)))
+ .collect::<Result<_>>()?;
+ self.scratch_rows.clear();
+ self.row_converter
+ .append(&mut self.scratch_rows, &ob_arrays)?;
+
+ let k = self.k;
+ let mut replacements: usize = 0;
+
+ // 4. Per-partition: bucket this batch's rows by distinct ob value
+ // (within-call accumulation), then merge each bucket into the
+ // partition state as a single `TieEntry`.
+ for (pk, indices) in groups.drain() {
+ let state =
+ self.states
+ .entry(pk)
+ .or_insert_with(|| DenseRankPartitionState {
+ groups: HashMap::new(),
+ max_key: None,
+ });
+
+ // Bucket by ob key. `ob_runs` is a reused scratch map (taken
+ // out and drained below) so its backing table is allocated
+ // once, not once per distinct partition key. `entry_ref` owns
+ // the key only on Vacant, so repeated rows of the same ob
+ // value don't re-allocate.
+ let mut runs = std::mem::take(&mut self.ob_runs);
+ runs.clear();
+ for &orig_idx in &indices {
+ let ob_row = self.scratch_rows.row(orig_idx as usize);
+ runs.entry_ref(ob_row.as_ref()).or_default().push(orig_idx);
+ }
+
+ for (ob_key, run_indices) in runs.drain() {
+ // Case A: ob already tracked — push this batch's run as a
+ // new `TieEntry` (one entry per contributing batch, exactly
+ // like RANK pushing one tie entry per batch).
+ if let Some(entries) = state.groups.get_mut(&ob_key) {
+ entries.push(TieEntry {
+ batch: batch.clone(),
+ row_indices: run_indices,
+ batch_bytes: input_batch_bytes,
+ });
+ continue;
+ }
+
+ // Case B: new ob, room available.
+ if state.groups.len() < k {
+ state.groups.insert(
+ ob_key,
+ vec![TieEntry {
+ batch: batch.clone(),
+ row_indices: run_indices,
+ batch_bytes: input_batch_bytes,
+ }],
+ );
+ // A new distinct key may raise the boundary; drop the
+ // cached max so it is recomputed on the next Case C.
+ state.max_key = None;
+ continue;
+ }
+
+ // Case C: new ob, at K distinct keys. The largest tracked
+ // ob value is the admission boundary; it is cached in
+ // `max_key` and recomputed via an O(K) scan only when the
+ // cache is stale, so the common "ob >= max → drop" path is
+ // O(1). Scoped so the immutable borrow ends before mutation.
+ if state.max_key.is_none() {
+ state.max_key = state.groups.keys().max().cloned();
+ }
+ let is_smaller = state
+ .max_key
+ .as_deref()
+ .map(|max_key| ob_key.as_slice() < max_key)
+ .expect("state.groups has k >= 1 keys");
+
+ if is_smaller {
+ // Evict the entire max-key group and drop the cached
+ // max (recomputed on the next Case C).
+ let evicted_key = state.max_key.take().expect("max key
present");
Review Comment:
The cache only makes rejected candidates O(1). Every accepted candidate
invalidates it, so input arriving from worst to best still scans all K keys for
each new value and remains O(N × K). The benchmarks use K=2 and do not exercise
this case. Could we use an ordered structure or add a large-K improving-order
benchmark?
--
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]