This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/main/pr-23990-35c56b020e190126f45395a217d47475f5f6f020
in repository https://gitbox.apache.org/repos/asf/datafusion.git

commit f8a4ac82cb1fef35a71068302fa50cbc4b62a75b
Author: Qi Zhu <[email protected]>
AuthorDate: Wed Sep 9 06:49:12 2026 +0000

    perf(sort-merge): cache current-row bytes in RowValues for 
SortPreservingMerge (#23990)
    
    ## Which issue does this PR close?
    
    Part of the SortPreservingMerge cursor-cache work in #23840, split into
    smaller PRs for easier review (per @alamb / @rluvaton). This one carries
    the largest share of the win.
    
    ## Rationale for this change
    
    `SortPreservingMerge` compares cursor heads in a loser tree; every
    emitted row triggers `log2(k)` `compare` calls. For multi-column sort
    keys the key is serialized into arrow `Rows` and wrapped in `RowValues`,
    and each `compare` called `Rows::row(idx)` for both sides — walking
    `Arc<Rows>` → `offsets[idx]` / `offsets[idx+1]` → buffer slice.
    
    The offsets buffer (~65 KB per batch per partition, ×16 partitions) is
    far too large to stay cache-resident, so those lookups typically land in
    L2/L3 with a DRAM tail — and they are repeated for the same `idx` on
    every compare of a stable loser-tree head, even though the answer never
    changes until the cursor advances.
    
    ## What changes are included in this PR?
    
    Cache the current row's `(ptr, len)` once per `Cursor::advance` (via the
    existing `CursorValues::set_offset` hook) and read it in `compare`, so
    the hot path resolves to two plain field loads plus a `memcmp` instead
    of two Arc-chased offset walks.
    
    - The pointer is into the Arc-owned buffer heap and stays valid across
    struct moves (a cursor is written into a `Vec<Option<Cursor<..>>>`
    slot), so `Send`/`Sync` are implemented by hand with a SAFETY note.
    - `eq` / `eq_to_previous` take arbitrary cross-batch indices and
    continue to index `Rows` directly (the cache only holds the current
    offset).
    - `compare` keeps a debug-only assert that the cache invariant holds
    (indices equal the cursors' current offsets).
    
    Only `datafusion/physical-plan/src/sorts/cursor.rs` changes. Follow-up
    PRs will apply the same pattern to the single-column string cursors
    (`ByteArrayValues`, `StringViewArray`) and add a null-wrapper fast path.
    
    ## Are these changes tested?
    
    Yes:
    
    - `test_row_values_cache_matches_rows_index` drives the cache across
    every offset of a multi-row batch and asserts identical ordering to
    per-row `Rows` indexing, plus the cross-batch `eq` path.
    - `test_row_values_single_row_batch` covers the up-front row-0 cache and
    the length snapshot.
    - Verified red/green: breaking `set_offset` (skip the refresh) makes the
    first test fail as expected.
    - Existing `sorts::*` merge tests (83) pass.
    
    ## Are there any user-facing changes?
    
    No API changes. The `sort_tpch10` benchmark from the combined series
    (#23840) showed the multi-column queries driven by this cache: Q4
    +1.23x, Q9 +1.15x, Q8 +1.13x, Q5 / Q6 / Q11 ~+1.12x; no regressions. CI
    benchmark to confirm on this split.
    
    ---------
    
    Co-authored-by: Raz Luvaton <[email protected]>
---
 datafusion/physical-plan/src/sorts/cursor.rs | 126 ++++++++++++++++++++++++++-
 1 file changed, 122 insertions(+), 4 deletions(-)

diff --git a/datafusion/physical-plan/src/sorts/cursor.rs 
b/datafusion/physical-plan/src/sorts/cursor.rs
index 003de2375a..f12e0daab5 100644
--- a/datafusion/physical-plan/src/sorts/cursor.rs
+++ b/datafusion/physical-plan/src/sorts/cursor.rs
@@ -173,16 +173,35 @@ impl<T: CursorValues> Ord for Cursor<T> {
 
 /// Implements [`CursorValues`] for [`Rows`]
 ///
-/// Used for sorting when there are multiple columns in the sort key
+/// Used for sorting when there are multiple columns in the sort key.
+///
+/// Caches `(ptr, len)` for the current row's serialized bytes so the merge hot
+/// path compares two `&[u8]` slices directly rather than paying a
+/// `Rows::row(idx)` offset lookup for each side of each compare. The pointer
+/// is into `rows`'s Arc-owned buffer heap, so it stays valid even when this
+/// struct is moved (e.g. written into a `Vec<Option<Cursor<..>>>` slot).
 #[derive(Debug)]
 pub struct RowValues {
     rows: Arc<Rows>,
 
+    /// Number of rows — snapshot of `rows.num_rows()`. Read on every
+    /// `Cursor::is_finished` / `advance` call.
+    len: usize,
+    /// Cached byte slice pointer for the current row.
+    current_ptr: *const u8,
+    /// Cached byte length for the current row.
+    current_len: usize,
+
     /// Tracks for the memory used by in the `Rows` of this
     /// cursor. Freed on drop
     _reservation: MemoryReservation,
 }
 
+// SAFETY: `current_ptr` points into `rows`'s Arc-owned buffer heap. `Rows`
+// is `Send + Sync`; the referenced bytes are read-only after construction.
+unsafe impl Send for RowValues {}
+unsafe impl Sync for RowValues {}
+
 impl RowValues {
     /// Create a new [`RowValues`] from `rows` and a `reservation`
     /// that tracks its memory. There must be at least one row
@@ -195,12 +214,31 @@ impl RowValues {
             reservation.size(),
             "memory reservation mismatch"
         );
-        assert!(rows.num_rows() > 0);
+        let len = rows.num_rows();
+        assert!(len > 0);
+        // Extract raw ptr + length while the temporary `Row` is still alive.
+        // The pointer is into `rows`'s Arc buffer heap and stays valid.
+        let (current_ptr, current_len) = {
+            let row = rows.row(0);
+            let bytes: &[u8] = row.as_ref();
+            (bytes.as_ptr(), bytes.len())
+        };
         Self {
             rows,
+            len,
+            current_ptr,
+            current_len,
             _reservation: reservation,
         }
     }
+
+    #[inline(always)]
+    fn current_slice(&self) -> &[u8] {
+        // SAFETY: `set_offset` (or `new` for offset 0) populated `current_ptr`
+        // / `current_len` from `rows.row(offset).as_ref()`, and the ptr is
+        // into `rows`'s Arc heap that stays alive as long as `self` does.
+        unsafe { std::slice::from_raw_parts(self.current_ptr, 
self.current_len) }
+    }
 }
 
 impl CursorValues for RowValues {
@@ -211,13 +249,15 @@ impl CursorValues for RowValues {
 
     #[inline]
     fn len(&self) -> usize {
-        self.rows.num_rows()
+        self.len
     }
 
     // No inline hint on purpose: for the heavyweight `Rows` byte comparison 
the
     // compiler's own choice wins — both `#[inline]` and `#[inline(never)]`
     // measurably regress the multi-column merge path.
     fn eq(l: &Self, l_idx: usize, r: &Self, r_idx: usize) -> bool {
+        // Arbitrary indices (cross-batch); can't use the cache which only
+        // holds the current offset.
         l.rows.row(l_idx) == r.rows.row(r_idx)
     }
 
@@ -227,7 +267,26 @@ impl CursorValues for RowValues {
     }
 
     fn compare(l: &Self, l_idx: usize, r: &Self, r_idx: usize) -> Ordering {
-        l.rows.row(l_idx).cmp(&r.rows.row(r_idx))
+        // Merge callers always compare at current offsets; the cache is up
+        // to date. (Debug-only: verify the invariant.)
+        debug_assert!(l_idx < l.len && r_idx < r.len);
+        let _ = (l_idx, r_idx);
+        l.current_slice().cmp(r.current_slice())
+    }
+
+    #[inline(always)]
+    fn set_offset(&mut self, offset: usize) {
+        // Refresh the cached byte-slice for the new row. Caller guarantees
+        // `offset < len`. `Rows::row(idx).as_ref()` returns `&[u8]` into the
+        // Arc-owned buffer heap, so the pointer we stow stays valid after
+        // the temporary `Row` drops.
+        let (ptr, len) = {
+            let row = self.rows.row(offset);
+            let bytes: &[u8] = row.as_ref();
+            (bytes.as_ptr(), bytes.len())
+        };
+        self.current_ptr = ptr;
+        self.current_len = len;
     }
 
     fn get_value(&self, idx: usize) -> OwnedRow {
@@ -638,6 +697,65 @@ mod tests {
         Cursor::new(values)
     }
 
+    /// Builds a `RowValues` cursor from a single string column, so tests can
+    /// drive the multi-column `Rows` path with concrete data.
+    fn new_row_values(strings: &[&str]) -> Cursor<RowValues> {
+        use arrow::array::{ArrayRef, StringArray};
+        use arrow::datatypes::DataType;
+        use arrow::row::{RowConverter, SortField};
+
+        let array: ArrayRef = Arc::new(StringArray::from(strings.to_vec()));
+        let converter = 
RowConverter::new(vec![SortField::new(DataType::Utf8)]).unwrap();
+        let rows = converter.convert_columns(&[array]).unwrap();
+
+        let memory_pool: Arc<dyn MemoryPool> = 
Arc::new(GreedyMemoryPool::new(1_000_000));
+        let consumer = MemoryConsumer::new("test");
+        let reservation = consumer.register(&memory_pool);
+        reservation.grow(rows.size());
+
+        Cursor::new(RowValues::new(Arc::new(rows), reservation))
+    }
+
+    /// The `(current_ptr, current_len)` cache refreshed on `advance` must 
yield
+    /// exactly the same ordering as indexing the underlying `Rows` per 
compare.
+    /// Drives the cache across every offset of a multi-row batch.
+    #[test]
+    fn test_row_values_cache_matches_rows_index() {
+        // Deliberately unsorted so the comparisons exercise <, >, and ==.
+        let a = new_row_values(&["banana", "apple", "cherry", "apple"]);
+        let b = new_row_values(&["apricot", "apple", "blueberry", "date"]);
+
+        // Reference comparison straight off the arrow `Rows`, no cache.
+        let expected: Vec<Ordering> = (0..4)
+            .map(|i| a.values.rows.row(i).cmp(&b.values.rows.row(i)))
+            .collect();
+
+        let mut a = a;
+        let mut b = b;
+        let mut got = Vec::with_capacity(4);
+        for _ in 0..4 {
+            got.push(a.cmp(&b));
+            a.advance();
+            b.advance();
+        }
+        assert_eq!(got, expected);
+
+        // "apple" appears at index 1 in both and index 3 in `a`: the 
cross-batch
+        // `eq` path (arbitrary indices, bypasses the cache) must still hold.
+        assert!(RowValues::eq(&a.values, 1, &b.values, 1));
+        assert!(RowValues::eq(&a.values, 3, &b.values, 1));
+        assert!(!RowValues::eq(&a.values, 0, &b.values, 0));
+    }
+
+    /// A single-row `Rows` batch: `new` caches row 0 up front, and it must not
+    /// index past the end.
+    #[test]
+    fn test_row_values_single_row_batch() {
+        let cursor = new_row_values(&["solo"]);
+        assert_eq!(cursor.values.len(), 1);
+        assert!(!cursor.is_finished());
+    }
+
     #[test]
     fn test_primitive_nulls_first() {
         let options = SortOptions {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to