kosiew commented on code in PR #25716:
URL: https://github.com/apache/datafusion/pull/25716#discussion_r4226426693


##########
datafusion/physical-plan/src/joins/hash_join/exec.rs:
##########
@@ -3144,6 +3149,133 @@ fn concat_build_batches(
     Ok(batch)
 }
 
+/// Deduplicates shared data buffer references in a [`GenericByteViewArray`] 
by pointer identity.
+///
+/// When multiple record batches that share underlying buffer allocations are 
concatenated,
+/// Arrow's `concat` kernel appends every batch's `data_buffers` list 
verbatim, resulting
+/// in N × K buffer references for N batches that share K allocations.
+///
+/// This function walks the buffer list, identifies duplicates by raw pointer 
address,
+/// and rewrites the 4-byte `buffer_index` inside each non-inline view (length 
> 12) to
+/// point into the deduplicated buffer vector. **No string bytes are copied.**
+///
+/// The fast path (0 or 1 data buffers, or no duplicates found) clones the 
array reference

Review Comment:
   The no-duplicates path still allocates temporary lookup structures when 
there are multiple buffers. Could we clarify that it avoids rewriting the views 
rather than claiming it performs no allocations?



##########
datafusion/physical-plan/src/joins/hash_join/exec.rs:
##########
@@ -3144,6 +3149,133 @@ fn concat_build_batches(
     Ok(batch)
 }
 
+/// Deduplicates shared data buffer references in a [`GenericByteViewArray`] 
by pointer identity.
+///
+/// When multiple record batches that share underlying buffer allocations are 
concatenated,
+/// Arrow's `concat` kernel appends every batch's `data_buffers` list 
verbatim, resulting
+/// in N × K buffer references for N batches that share K allocations.
+///
+/// This function walks the buffer list, identifies duplicates by raw pointer 
address,

Review Comment:
   Could we clarify that deduplication uses identical (pointer, length) pairs 
rather than allocation identity alone? Different ranges within the same 
allocation must remain separate.



##########
datafusion/physical-plan/src/joins/hash_join/exec.rs:
##########
@@ -11200,4 +11664,355 @@ mod tests {
         assert!(join.set_dynamic_filter(df).is_err());
         Ok(())
     }
+
+    // -----------------------------------------------------------------------
+    // Unit tests for deduplicate_view_array_buffers /
+    //                deduplicate_record_batch_view_buffers
+    // -----------------------------------------------------------------------
+
+    /// Fast path: an array with a single data buffer must be returned as-is
+    /// (pointer-equal clone, zero allocations).
+    #[test]
+    fn test_dedup_view_array_single_buffer_is_noop() {
+        let array = StringViewArray::from(vec![
+            "hello world long string!",
+            "another long value!",
+        ]);
+        let result = deduplicate_view_array_buffers(&array);
+        // Same number of data buffers — nothing removed.
+        assert_eq!(result.data_buffers().len(), array.data_buffers().len());
+    }
+
+    /// Fast path: an array with zero data buffers (all inline values) must be 
returned as-is.
+    #[test]
+    fn test_dedup_view_array_zero_buffers_is_noop() {
+        let array = StringViewArray::from(vec!["inline_only"]);
+        assert_eq!(array.data_buffers().len(), 0);
+        let result = deduplicate_view_array_buffers(&array);
+        assert_eq!(result.data_buffers().len(), 0);
+        assert_eq!(result.value(0), "inline_only");
+    }
+
+    /// Fast path: multiple distinct buffers (no duplicates) → returned as-is.
+    #[test]
+    fn test_dedup_view_array_no_duplicates_is_noop() {
+        // Build two separate StringViewArrays so their buffers are distinct.
+        let a = StringViewArray::from(vec!["first long string value here"]);
+        let b = StringViewArray::from(vec!["second long string value here"]);
+        // Concatenate to get an array with two *different* buffers.
+        let combined = arrow::compute::concat(&[&a as _, &b as _]).unwrap();
+        let combined = 
combined.as_any().downcast_ref::<StringViewArray>().unwrap();
+        let n_before = combined.data_buffers().len();
+        let result = deduplicate_view_array_buffers(combined);
+        // Still the same number of unique buffers — nothing deduplicated.
+        assert_eq!(result.data_buffers().len(), n_before);
+    }
+
+    /// Deduplication path: concatenating an array with itself produces N refs 
to
+    /// the same K buffers; deduplicate_view_array_buffers collapses them to K.
+    #[test]
+    fn test_dedup_view_array_removes_duplicate_buffers() {
+        let base = StringViewArray::from(vec!["this is a long string value 
abc"]);
+        // concat gives us 2 references to the same underlying buffer.
+        let doubled = arrow::compute::concat(&[&base as _, &base as 
_]).unwrap();
+        let doubled = 
doubled.as_any().downcast_ref::<StringViewArray>().unwrap();
+        assert_eq!(doubled.data_buffers().len(), 2);
+        let deduped = deduplicate_view_array_buffers(doubled);
+        // After deduplication only 1 unique buffer should remain.
+        assert_eq!(deduped.data_buffers().len(), 1);
+        // Values must be preserved.
+        assert_eq!(deduped.value(0), "this is a long string value abc");
+        assert_eq!(deduped.value(1), "this is a long string value abc");
+    }
+
+    /// Inline-string path: strings with length ≤ 12 are stored inline in the
+    /// view descriptor and carry no buffer index; they must survive 
deduplication.
+    #[test]
+    fn test_dedup_view_array_inline_strings_preserved() {
+        // "hi" is 2 bytes — well within the 12-byte inline threshold.
+        let base = StringViewArray::from(vec!["hi", "short"]);
+        let doubled = arrow::compute::concat(&[&base as _, &base as 
_]).unwrap();
+        let doubled = 
doubled.as_any().downcast_ref::<StringViewArray>().unwrap();
+        let deduped = deduplicate_view_array_buffers(doubled);
+        assert_eq!(deduped.value(0), "hi");
+        assert_eq!(deduped.value(1), "short");
+        assert_eq!(deduped.value(2), "hi");
+        assert_eq!(deduped.value(3), "short");
+    }
+
+    /// Mixed inline, long strings, and nulls: ensures the view rewriting loop 
handles
+    /// non-inline view remapping, skips inline views (length <= 12), and 
preserves the null buffer.
+    #[test]
+    fn test_dedup_view_array_mixed_inline_long_and_nulls() {
+        let array = StringViewArray::from(vec![
+            Some("this is a long string that exceeds inline size 12"),
+            Some("inline"),
+            None,
+        ]);
+        let doubled = arrow::compute::concat(&[&array as _, &array as 
_]).unwrap();
+        let doubled = 
doubled.as_any().downcast_ref::<StringViewArray>().unwrap();
+        assert_eq!(doubled.data_buffers().len(), 2);
+        let deduped = deduplicate_view_array_buffers(doubled);
+        assert_eq!(deduped.data_buffers().len(), 1);
+        assert_eq!(
+            deduped.value(0),
+            "this is a long string that exceeds inline size 12"
+        );
+        assert_eq!(deduped.value(1), "inline");
+        assert!(deduped.is_null(2));
+        assert_eq!(
+            deduped.value(3),
+            "this is a long string that exceeds inline size 12"
+        );
+        assert_eq!(deduped.value(4), "inline");
+        assert!(deduped.is_null(5));
+    }
+
+    /// Direct BinaryView path: exercises deduplicate_view_array_buffers 
directly on BinaryViewArray
+    /// with mixed long values, inline values, and nulls.
+    #[test]
+    fn test_dedup_view_array_binary_view_direct() {
+        let base = BinaryViewArray::from_iter(vec![
+            Some(&b"this is definitely longer than 12 bytes"[..]),
+            Some(&b"short"[..]),
+            None,
+        ]);
+        let doubled = arrow::compute::concat(&[&base as _, &base as 
_]).unwrap();
+        let doubled = 
doubled.as_any().downcast_ref::<BinaryViewArray>().unwrap();
+        assert_eq!(doubled.data_buffers().len(), 2);
+        let deduped = deduplicate_view_array_buffers(doubled);
+        assert_eq!(deduped.data_buffers().len(), 1);
+        assert_eq!(deduped.value(0), b"this is definitely longer than 12 
bytes");
+        assert_eq!(deduped.value(1), b"short");
+        assert!(deduped.is_null(2));
+        assert_eq!(deduped.value(3), b"this is definitely longer than 12 
bytes");
+        assert_eq!(deduped.value(4), b"short");
+        assert!(deduped.is_null(5));
+    }
+
+    /// BinaryView path: `deduplicate_record_batch_view_buffers` must also
+    /// deduplicate `BinaryView` columns (exercises the BinaryView arm).
+    #[test]
+    fn test_dedup_record_batch_binary_view_column() {
+        let schema = Arc::new(Schema::new(vec![Field::new(
+            "data",
+            DataType::BinaryView,
+            false,
+        )]));
+        let base = BinaryViewArray::from_iter_values(vec![
+            b"this is definitely longer than 12 bytes",
+        ]);
+        let doubled = arrow::compute::concat(&[&base as _, &base as 
_]).unwrap();
+        let batch = RecordBatch::try_new(Arc::clone(&schema), 
vec![doubled]).unwrap();
+        assert_eq!(
+            batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<BinaryViewArray>()
+                .unwrap()
+                .data_buffers()
+                .len(),
+            2
+        );
+        let deduped = deduplicate_record_batch_view_buffers(&batch);
+        let result = deduped
+            .column(0)
+            .as_any()
+            .downcast_ref::<BinaryViewArray>()
+            .unwrap();
+        assert_eq!(result.data_buffers().len(), 1);
+    }
+
+    /// Mixed columns path: `deduplicate_record_batch_view_buffers` must 
deduplicate
+    /// both Utf8View and BinaryView columns while passing through non-view 
columns (Int32, Utf8).
+    #[test]
+    fn test_dedup_record_batch_mixed_view_and_non_view_columns() {
+        let schema = Arc::new(Schema::new(vec![
+            Field::new("s", DataType::Utf8View, true),
+            Field::new("b", DataType::BinaryView, true),
+            Field::new("n", DataType::Int32, false),
+            Field::new("str", DataType::Utf8, false),
+        ]));
+        let str_base = StringViewArray::from(vec!["long string value here 
abc"]);
+        let bin_base =
+            BinaryViewArray::from_iter_values(vec![b"long binary value here 
abc"]);
+        let str_doubled =
+            arrow::compute::concat(&[&str_base as _, &str_base as _]).unwrap();
+        let bin_doubled =
+            arrow::compute::concat(&[&bin_base as _, &bin_base as _]).unwrap();
+        let int_col: ArrayRef = Arc::new(Int32Array::from(vec![42, 43]));
+        let utf8_col: ArrayRef = Arc::new(StringArray::from(vec!["foo", 
"bar"]));
+
+        let batch = RecordBatch::try_new(
+            Arc::clone(&schema),
+            vec![
+                str_doubled,
+                bin_doubled,
+                Arc::clone(&int_col),
+                Arc::clone(&utf8_col),
+            ],
+        )
+        .unwrap();
+
+        let deduped = deduplicate_record_batch_view_buffers(&batch);
+        let deduped_str = deduped
+            .column(0)
+            .as_any()
+            .downcast_ref::<StringViewArray>()
+            .unwrap();
+        let deduped_bin = deduped
+            .column(1)
+            .as_any()
+            .downcast_ref::<BinaryViewArray>()
+            .unwrap();
+
+        assert_eq!(deduped_str.data_buffers().len(), 1);
+        assert_eq!(deduped_bin.data_buffers().len(), 1);
+        // Non-view columns passed through unmodified
+        assert!(Arc::ptr_eq(deduped.column(2), &int_col));
+        assert!(Arc::ptr_eq(deduped.column(3), &utf8_col));
+    }
+
+    /// Fast path: a batch with no Utf8View or BinaryView columns must be
+    /// returned as a cheap clone with no allocations.
+    #[test]
+    fn test_dedup_record_batch_no_view_columns_is_noop() {
+        let schema = Arc::new(Schema::new(vec![Field::new("n", 
DataType::Int32, false)]));
+        let batch = RecordBatch::try_new(
+            Arc::clone(&schema),
+            vec![Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef],
+        )
+        .unwrap();
+        // Should succeed without error and return the same data.
+        let result = deduplicate_record_batch_view_buffers(&batch);
+        assert_eq!(result.num_rows(), 3);
+    }
+
+    /// Reverse concatenation order: exercises concat_build_batches with 
reverse=true
+    /// on view columns to verify both reverse order and buffer deduplication 
work together.
+    #[test]
+    fn test_concat_build_batches_reverse_order_deduplication() {
+        use arrow::array::StringViewBuilder;
+
+        let schema = Arc::new(Schema::new(vec![Field::new(
+            "s",
+            DataType::Utf8View,
+            false,
+        )]));
+
+        let mut builder = StringViewBuilder::new();
+        builder.append_value("batch1: long string value exceeding twelve 
bytes");
+        let batch1 = RecordBatch::try_new(
+            Arc::clone(&schema),
+            vec![Arc::new(builder.finish()) as ArrayRef],
+        )
+        .expect("valid batch");
+
+        let mut builder = StringViewBuilder::new();
+        builder.append_value("batch2: long string value exceeding twelve 
bytes");
+        let batch2 = RecordBatch::try_new(
+            Arc::clone(&schema),
+            vec![Arc::new(builder.finish()) as ArrayRef],
+        )
+        .expect("valid batch");
+
+        let metrics = BuildProbeJoinMetrics::new(0, 
&ExecutionPlanMetricsSet::new());
+        let pool: Arc<dyn MemoryPool> = 
Arc::new(UnboundedMemoryPool::default());
+        // 4 batches with two duplicate references each — raw concat yields 4 
buffer handles
+        let batches = vec![batch1.clone(), batch2.clone(), batch1, batch2];
+        let (mut reservation, inputs_reserved) =
+            reserve_inputs(&batches, &pool).expect("reserve");
+
+        let batch = concat_build_batches(
+            &schema,
+            batches,
+            true,
+            inputs_reserved,
+            &mut reservation,
+            &metrics,
+        )
+        .expect("concat");
+
+        let view_arr = batch
+            .column(0)
+            .as_any()
+            .downcast_ref::<StringViewArray>()
+            .unwrap();
+
+        // After deduplication 4 handles collapse to 2 unique buffers.
+        assert_eq!(view_arr.data_buffers().len(), 2);
+        assert_eq!(batch.num_rows(), 4);
+        // reverse=true: batch2, batch1, batch2, batch1
+        assert_eq!(
+            view_arr.value(0),
+            "batch2: long string value exceeding twelve bytes"
+        );
+        assert_eq!(
+            view_arr.value(1),
+            "batch1: long string value exceeding twelve bytes"
+        );
+        assert_eq!(
+            view_arr.value(2),
+            "batch2: long string value exceeding twelve bytes"
+        );
+        assert_eq!(
+            view_arr.value(3),
+            "batch1: long string value exceeding twelve bytes"
+        );
+    }
+
+    /// Exercises the `retained > held` growth path in concat_build_batches.
+    ///
+    /// When a single batch is provided `copy_size` is 0, so we only 
pre-reserve the
+    /// input size. If the returned batch happens to occupy more memory than 
what we
+    /// reserved (e.g., due to deduplication creating new buffers or tracking 
overhead
+    /// differences), the function must call `reservation.try_grow(retained - 
held)`.
+    /// We trigger this by passing `inputs_reserved = 0` directly so that 
`held == 0`
+    /// and any non-empty batch forces the grow branch.
+    #[test]
+    fn test_concat_build_batches_grow_branch() {
+        use arrow::array::StringViewBuilder;
+
+        let schema = Arc::new(Schema::new(vec![Field::new(
+            "s",
+            DataType::Utf8View,
+            false,
+        )]));
+
+        let mut builder = StringViewBuilder::new();
+        builder.append_value("long string that exceeds the twelve byte inline 
threshold");
+        let batch = RecordBatch::try_new(
+            Arc::clone(&schema),
+            vec![Arc::new(builder.finish()) as ArrayRef],
+        )
+        .expect("valid batch");
+
+        let metrics = BuildProbeJoinMetrics::new(0, 
&ExecutionPlanMetricsSet::new());
+        let pool: Arc<dyn MemoryPool> = 
Arc::new(UnboundedMemoryPool::default());
+        let mut reservation = MemoryConsumer::new("test").register(&pool);
+
+        // Pass inputs_reserved = 0 so that held == 0 + copy_size == 0 == 0,
+        // while retained == get_record_batch_memory_size(&batch) > 0 — forcing
+        // the `retained > held` branch to call reservation.try_grow.
+        let result = concat_build_batches(
+            &schema,
+            vec![batch],
+            false,
+            0, // inputs_reserved deliberately zero
+            &mut reservation,
+            &metrics,
+        )
+        .expect("concat");
+
+        let view_arr = result
+            .column(0)
+            .as_any()
+            .downcast_ref::<StringViewArray>()
+            .unwrap();
+        assert_eq!(result.num_rows(), 1);
+        assert_eq!(
+            view_arr.value(0),
+            "long string that exceeds the twelve byte inline threshold"
+        );
+    }

Review Comment:
   Could we strengthen this test by checking the final reservation size? For 
example, if the existing accounting contract requires the reservation to equal 
the retained batch size, assert_eq!(reservation.size(), 
get_record_batch_memory_size(&result)) would verify the accounting result 
rather than just the returned values.



##########
datafusion/physical-plan/src/joins/hash_join/exec.rs:
##########
@@ -3144,6 +3149,133 @@ fn concat_build_batches(
     Ok(batch)
 }
 
+/// Deduplicates shared data buffer references in a [`GenericByteViewArray`] 
by pointer identity.
+///
+/// When multiple record batches that share underlying buffer allocations are 
concatenated,
+/// Arrow's `concat` kernel appends every batch's `data_buffers` list 
verbatim, resulting
+/// in N × K buffer references for N batches that share K allocations.
+///
+/// This function walks the buffer list, identifies duplicates by raw pointer 
address,
+/// and rewrites the 4-byte `buffer_index` inside each non-inline view (length 
> 12) to
+/// point into the deduplicated buffer vector. **No string bytes are copied.**
+///
+/// The fast path (0 or 1 data buffers, or no duplicates found) clones the 
array reference
+/// with no allocations.
+fn deduplicate_view_array_buffers<T: ByteViewType>(
+    array: &GenericByteViewArray<T>,
+) -> GenericByteViewArray<T> {
+    let data_buffers = array.data_buffers();
+    if data_buffers.len() <= 1 {
+        return array.clone();
+    }
+
+    // Use the raw buffer address and length as the deduplication key. Casting 
to usize is the
+    // idiomatic way to use pointer values as HashMap keys on stable Rust.
+    let mut unique_buffers: Vec<arrow::buffer::Buffer> =
+        Vec::with_capacity(data_buffers.len());
+    let mut pointer_map: HashMap<(usize, usize), u32> =
+        HashMap::with_capacity(data_buffers.len());
+    let mut index_remap: Vec<u32> = Vec::with_capacity(data_buffers.len());
+    let mut has_duplicates = false;
+
+    for buf in data_buffers.iter() {
+        let key = (buf.as_ptr() as usize, buf.len());
+        if let Some(&new_idx) = pointer_map.get(&key) {
+            index_remap.push(new_idx);
+            has_duplicates = true;
+        } else {
+            let new_idx = unique_buffers.len() as u32;
+            pointer_map.insert(key, new_idx);
+            unique_buffers.push(buf.clone());
+            index_remap.push(new_idx);
+        }
+    }
+
+    if !has_duplicates {
+        return array.clone();
+    }
+
+    // Rewrite the buffer_index field in the 128-bit view descriptor for every
+    // non-inline value. Inline values (length <= 12) embed the payload inside
+    // the descriptor itself and carry no buffer index, so they are left as-is.
+    let views = array.views();
+    let mut new_views: Vec<u128> = Vec::with_capacity(views.len());
+    for &v in views.iter() {
+        let mut view = ByteView::from(v);
+        if view.length > 12 {
+            view.buffer_index = index_remap[view.buffer_index as usize];
+        }
+        new_views.push(view.as_u128());
+    }
+
+    let new_views_buffer = ScalarBuffer::from(new_views);
+    let nulls = array.nulls().cloned();
+
+    // SAFETY: `new_views_buffer` contains only valid 128-bit view descriptors
+    // derived from the source array. Each non-inline view's `buffer_index` has
+    // been remapped to point at the logically equivalent deduplicated buffer 
in
+    // `unique_buffers`, preserving the original byte offsets and lengths.
+    unsafe {
+        GenericByteViewArray::<T>::new_unchecked(
+            new_views_buffer,
+            unique_buffers.into(),
+            nulls,
+        )
+    }
+}
+
+/// Deduplicates shared data buffer references across all `Utf8View` and 
`BinaryView`
+/// columns in a [`RecordBatch`], returning a new batch whose view arrays hold 
at most
+/// as many buffer references as there are distinct underlying allocations.
+///
+/// Columns of other types are passed through unchanged. If the batch contains 
no view
+/// columns this function returns a cheap clone of the batch reference.
+/// Returns a new [`RecordBatch`] with deduplicated view buffer references.
+///
+/// This function is infallible: we reconstruct the batch using the original
+/// schema and the same set of columns (same lengths, same types). 
`RecordBatch::try_new`
+/// only fails when column lengths or schema mismatches occur — neither can 
happen here
+/// since we only replace view columns with logically equivalent deduplicated 
versions.

Review Comment:
   RecordBatch::try_new also checks for nulls in non-nullable fields. Could we 
update this comment to explain that the original valid batch's schema, column 
count, lengths, types, and null validity are all preserved?



-- 
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]

Reply via email to