kosiew commented on code in PR #25716:
URL: https://github.com/apache/datafusion/pull/25716#discussion_r4163816472
##########
datafusion/physical-plan/src/joins/hash_join/exec.rs:
##########
@@ -2986,6 +2992,121 @@ 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 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, 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 addr = buf.as_ptr() as usize;
Review Comment:
Using only `as_ptr()` is not enough to identify an Arrow `Buffer` range
because two valid slices can start at the same address but have different
lengths. Please deduplicate only identical ranges, for example by using pointer
plus length, or otherwise prove the retained buffer covers every remapped view,
and add a regression where a shorter slice is seen before the longer buffer and
the concatenated values are verified.
##########
datafusion/physical-plan/src/joins/hash_join/exec.rs:
##########
@@ -7607,6 +7728,65 @@ mod tests {
Ok(())
}
+ #[test]
+ fn concat_build_batches_deduplicates_view_buffers() -> Result<()> {
Review Comment:
Could you add coverage for `BinaryViewArray` as well? The implementation has
separate `BinaryView` dispatch and downcast logic, but the current regression
only exercises `StringViewArray`.
--
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]