Skip to content

fix: deduplicate StringView/BinaryView buffer refs in CollectLeft Has... - #25716

Open
mohitgurav20 wants to merge 19 commits into
apache:mainfrom
mohitgurav20:fix/dedup-stringview-buffers-collectleft-joins
Open

mohitgurav20 wants to merge 19 commits into
apache:mainfrom
mohitgurav20:fix/dedup-stringview-buffers-collectleft-joins

Conversation

@mohitgurav20

@mohitgurav20 mohitgurav20 commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #25712

Rationale for this change

Queries with multiple chained CollectLeft hash joins on StringView or BinaryView
columns can run 10x slower and use 7x more memory than expected. This happens even when
the underlying string data is tiny.

Why it happens: When Arrow's concat kernel combines batches that share the same
buffer allocations, it appends each batch's buffer list verbatim without checking for
duplicates. If N batches each hold a reference to the same K buffers, the result carries
N × K references — all pointing to the same memory. Every subsequent CollectLeft join
multiplies the count again. On TPC-DS SF1 Q64 with pushdown_filters = true we measured
a single column growing from 1 to 1,587,600 buffer references while the distinct
allocation count stayed at 2.

The two user-visible symptoms are:

  1. Query slowdown — RecordBatchMemoryCounter::count_buffer_memory_size walks every
    reference on each poll. With millions of duplicate refs, this becomes the dominant
    CPU cost.
  2. High memory use — the Arc<Buffer> reference vectors themselves occupy memory,
    and downstream take operations carry the full bloated buffer list into the next join.

On TPC-DS SF1 Q64:

pushdown_filters = false pushdown_filters = true
Query time (before) 0.75–1.4 s 4.0–10.6 s
Query time (after) 0.75–1.4 s 0.75–0.85 s
Peak RSS (before) ~1.2 GB ~8.5 GB
Peak RSS (after) ~1.2 GB ~1.2 GB

What changes are included in this PR?

A single new step is inserted immediately after the build-side concat_batches call in
concat_build_batches (hash_join/exec.rs):

  1. deduplicate_view_array_buffers<T> — for a single GenericByteViewArray, walks
    the data_buffers list, identifies duplicates by raw pointer address, and rewrites
    the 4-byte buffer_index inside each non-inline view descriptor to point into the
    deduplicated buffer vector. No string bytes are copied. The fast path (≤ 1 buffer, or
    no duplicates) returns a cheap clone() of the array reference.

  2. deduplicate_record_batch_view_buffers — calls the above for every Utf8View
    and BinaryView column in a RecordBatch. Columns of other types are passed through
    with Arc::clone.

The deduplication runs once, on the concatenated build batch, before it is handed to the
join probe loop. All subsequent take operations on the build side then start from a
clean buffer list.

What is the testing strategy for this PR?

A new unit test concat_build_batches_deduplicates_view_buffers (in the tests module
of hash_join/exec.rs) constructs three RecordBatches that share the same
StringViewArray allocation, concatenates them through concat_build_batches, and
asserts that the resulting column holds exactly one buffer reference instead of three.
It also asserts the expected row count to confirm no data was dropped or duplicated.

The existing hash-join test suite continues to pass without modification, confirming the
change is a no-op for non-view columns and for view columns that already have unique
buffers.

Are there any user-facing changes?

No API changes. The fix is entirely internal to concat_build_batches. Users running
queries with chained CollectLeft joins on StringView or BinaryView columns will
see lower memory usage and faster query times without any configuration change.

@github-actions github-actions Bot added physical-plan Changes to the physical-plan crate auto detected api change Auto detected API change labels Sep 24, 2026
@mohitgurav20

Copy link
Copy Markdown
Contributor Author

Hi @alamb @kosiew @asolimando,

I've formatted the code and rebased the branch onto latest main. This resolves the cargo-semver-checks false-positive (StatisticsExec::with_partition_statistics) that was triggered by the base branch being out-of-date.

All CI checks are now running against latest main. Ready for review when you have a moment!

@alamb

alamb commented Sep 27, 2026

Copy link
Copy Markdown
Contributor

Can you please update the description of this PR to be more readable

Also, perhaps @adriangb would be able to review it as he filed #25712

@github-actions github-actions Bot removed the auto detected api change Auto detected API change label Sep 27, 2026
@adriangb

Copy link
Copy Markdown
Contributor

Can you please update the description of this PR to be more readable

Highly agree, this is hard to read.

@mohitgurav20 a prompt I find very useful:

- Write the description focusing on the impact / experience for users (for example a SQL MRE), don't focus on implementation details.
- The target reader is a maintainer of DataFusion that is familiar with the space in general but may not be aware of this particular area / piece of the codebase.
- Define any terms you may be introducing or that would be unfamiliar to the reader above.
- Write using Simplified Technical English.
- Consider using tables, diagrams, example code snippets, lists or other forms of communication over prose where possible.

@mohitgurav20

Copy link
Copy Markdown
Contributor Author

Hey all,

I’ve fully updated the PR description following the suggested prompt to make it much easier to read (focusing on the user impact, adding the SQL MRE, and removing the heavy implementation jargon).

I also went through the CI logs and fixed the compilation issue causing the 5 checks to fail (it just needed a quick .into() conversion from Vec to Arc<[Buffer]> for the new_unchecked initialization).

Everything should build properly now. Let me know what you think!

@codecov-commenter

codecov-commenter commented Sep 27, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 82.74%. Comparing base (6bbd73a) to head (cd71ceb).
⚠️ Report is 3 commits behind head on main.

Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25716      +/-   ##
==========================================
+ Coverage   82.72%   82.74%   +0.01%     
==========================================
  Files        1147     1147              
  Lines      448093   448589     +496     
  Branches   448093   448589     +496     
==========================================
+ Hits       370689   371179     +490     
- Misses      54892    54896       +4     
- Partials    22512    22514       +2     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@mohitgurav20

Copy link
Copy Markdown
Contributor Author

Hey all, just a quick update!

I've resolved the failing CI checks. The formatting has been aligned using cargo fmt, and I resolved the strict Clippy warning by explicitly using Arc::clone(&schema) instead of relying on .clone().

Everything is green now and this should be fully ready for review! Let me know if you need any further tweaks.

@kosiew

kosiew commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

@mohitgurav20
Please amend the PR description as per .github/pull_request_template.md

@mohitgurav20

Copy link
Copy Markdown
Contributor Author

Updated the PR description to follow the template — it now clearly covers the rationale, what changed, testing, and user impact. Ready for review when you get a chance, @kosiew @alamb @adriangb.

@kosiew kosiew left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@mohitgurav20,

Thanks for working on this. I found one correctness issue in the buffer deduplication that needs to be addressed before this can merge. I also left a suggestion to add BinaryView coverage.

let mut has_duplicates = false;

for buf in data_buffers.iter() {
let addr = buf.as_ptr() as usize;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

}

#[test]
fn concat_build_batches_deduplicates_view_buffers() -> Result<()> {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could you add coverage for BinaryViewArray as well? The implementation has separate BinaryView dispatch and downcast logic, but the current regression only exercises StringViewArray.

@mohitgurav20
mohitgurav20 force-pushed the fix/dedup-stringview-buffers-collectleft-joins branch from f3bf727 to 124082a Compare October 3, 2026 09:42
@mohitgurav20

Copy link
Copy Markdown
Contributor Author

@kosiew Good catch on the slicing edge case—using just the pointer address definitely wasn't safe enough there.

I've just pushed a fix for this. The deduplication map now uses a (pointer_address, length) tuple as the key to properly identify exact buffer ranges.

I also added the test coverage you requested:

A new test specifically for BinaryViewArray deduplication.
A regression test that verifies the slicing behavior to ensure the retained buffer correctly covers the remapped views without collisions.
Everything is passing on my end. Let me know how it looks to you now!

@mohitgurav20
mohitgurav20 requested a review from kosiew October 3, 2026 10:31

@kosiew kosiew left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@mohitgurav20,

Thanks for the follow-up. The pointer-and-length key fixes the correctness issue, and the BinaryView coverage looks good. There is still one regression-test gap that needs to be addressed.

Arc::new(Schema::new(vec![Field::new("s", DataType::Utf8View, true)]));

// Create a shorter slice that starts at the same address
let short_slice = base_array.slice(0, 1);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test does not actually cover buffers with the same pointer but different lengths. GenericByteViewArray::slice only slices the views, so please construct two arrays using differently sized Buffer::slice_with_length(0, ...) ranges with the shorter one first, then assert that values from the longer range still survive concatenation.

@mohitgurav20
mohitgurav20 force-pushed the fix/dedup-stringview-buffers-collectleft-joins branch 2 times, most recently from ce9ad51 to 0fb7157 Compare October 5, 2026 16:52
@github-actions github-actions Bot added the common Related to common crate label Oct 5, 2026
@mohitgurav20
mohitgurav20 force-pushed the fix/dedup-stringview-buffers-collectleft-joins branch from 0fb7157 to 9647de3 Compare October 5, 2026 16:53
@mohitgurav20
mohitgurav20 force-pushed the fix/dedup-stringview-buffers-collectleft-joins branch from 9647de3 to 3995a49 Compare October 5, 2026 18:29
Comment thread .gitignore

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@mohitgurav20
Can you exclude this change from this PR?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

"Done! Removed the .gitignore change — it was accidentally corrupted to UTF-16 by cargo fmt on Windows. Restored it to match main exactly."

@github-actions github-actions Bot removed the common Related to common crate label Oct 6, 2026
Address Codecov partial/missing coverage gaps identified in PR review:

- test_dedup_view_array_zero_buffers_is_noop: covers the early-return path
  in deduplicate_view_array_buffers when data_buffers is empty (all inline).

- test_dedup_view_array_mixed_inline_long_and_nulls: exercises the view
  rewriting loop with a mix of long strings (>12 bytes, non-inline),
  short inline strings, and null entries — covering the inline-skip branch
  and null buffer preservation together.

- test_dedup_view_array_binary_view_direct: validates the same deduplication
  logic on BinaryViewArray directly (distinct from the Utf8View path).

- test_dedup_record_batch_mixed_view_and_non_view_columns: ensures
  deduplicate_record_batch_view_buffers deduplicates both Utf8View and
  BinaryView columns while Arc-cloning non-view columns (Int32, Utf8)
  unchanged — covering the _ => Arc::clone(col) arm.

- test_concat_build_batches_reverse_order_deduplication: calls
  concat_build_batches with reverse=true and 4 batches (two pairs of
  duplicates), verifying reversed row order and that 4 raw buffer handles
  collapse to 2 unique buffers after deduplication.

Enhanced concat_build_batches_deduplicates_view_buffers and
concat_build_batches_deduplicates_binary_view_buffers to include inline
values, null entries, and non-view (Int32) columns. Expanded
concat_build_batches_deduplicates_slices_regression from 2 to 4 batches
so has_duplicates is true and the deduplication code path actually runs.
… path

Two targeted fixes to push patch coverage toward 100%:

1. Remove uncovered error branches from test_concat_build_batches_reverse_order_deduplication.
   The previous version returned Result<()> and used ? on RecordBatch::try_new and
   concat_build_batches. Each ? generates two LLVM branches -- Ok (taken) and Err
   (never taken in tests) -- which Codecov reports as partial lines. Converting to
   () + .expect() collapses each call to a single branch, eliminating all partials.

2. Add test_concat_build_batches_grow_branch to exercise the retained > held path
   in concat_build_batches (line 2997). Prior tests always produced batches where
   deduplication shrank the retained size, so only the else/shrink branch at line
   3001 was ever hit. The new test passes inputs_reserved = 0 directly, making
   held == 0, so any non-empty batch forces the try_grow call on line 2998.
…erage gaps

Codecov reported multiple partial missing lines due to the Err branches of ?
operators never being taken.

By changing deduplicate_record_batch_view_buffers to return RecordBatch directly
instead of Result<RecordBatch> (since RecordBatch::try_new with the same schema
and matching column lengths cannot fail), we eliminate the ? operator at the
call site in concat_build_batches. We also remove the unneeded .unwrap()
calls from the test assertions.
@mohitgurav20

Copy link
Copy Markdown
Contributor Author

Thanks for the detailed review and patience, @kosiew!

I've pushed a final set of commits that fully resolves the Codecov coverage gaps.

The persistent "partial" coverage was actually due to LLVM branch instrumentation on ? operators in both the production code and tests (where the Err branches were never hit).

To fix this, I made deduplicate_record_batch_view_buffers infallible (since it mathematically cannot fail with matching schemas/column lengths), swapped test assertions to .expect(), and added targeted unit tests for the remaining edge cases (like the try_grow path).

Everything is green locally. Let me know if we're good to merge!

@mohitgurav20
mohitgurav20 requested a review from kosiew October 6, 2026 12:53
@kosiew

kosiew commented Oct 7, 2026

Copy link
Copy Markdown
Contributor

@mohitgurav20
Thanks for your patience and persistence with this.
I'll review this after you resolve the merge conflicts.

Upstream added BooleanArray and NullBuffer imports. Our branch added
BinaryViewArray, ByteView, GenericByteViewArray, StringViewArray, and
ScalarBuffer for the StringView buffer dedup feature. Combined both
import sets; no logic changes.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

physical-plan Changes to the physical-plan crate

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Chained CollectLeft hash joins multiply StringView buffer references (TPC-DS Q64: 8.5 GB, ~10x slower with pushdown_filters)

5 participants