Skip to content

perf: retain batched build storage for large hash joins - #25889

Draft
sunchao wants to merge 6 commits into
apache:mainfrom
sunchao:codex/upstream-batched-hash-join
Draft

sunchao wants to merge 6 commits into
apache:mainfrom
sunchao:codex/upstream-batched-hash-join

Conversation

@sunchao

@sunchao sunchao commented Sep 30, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Part of #23031. Related to #23032; this is a hash-join-specific alternative for comparing the implementation and performance tradeoffs, not a replacement for its nested-loop or piecewise-merge join work.

Rationale for this change

Hash joins currently concatenate the entire build relation into one Arrow batch. For a large build this requires admitting a second payload allocation while the input batches are still retained, and combines individually valid variable-width arrays into one offset-limited array.

This PR keeps the small-build path and can retain large builds as batches. The scope includes the full build/probe/output path and a reproducible benchmark suite, rather than only exposing a storage helper.

What changes are included in this PR?

  • Select compact versus batched storage using both retained memory and estimated logical copy size, keeping tiny slices of large parents compact.
  • Coalesce metadata-heavy independent flat inputs within 8 MiB / 8,192-row targets, preserve already substantial batches, account for copies before allocation, and release unused backing buffers and empty-batch charges.
  • Keep batch-local evaluated keys and a sparse logical-row directory; compare candidates without concatenating the keys or payload.
  • Reserve the actual retained metadata capacities and share probe preprocessing; keep only one source comparator alive at a time, including Arrow's dictionary-null scratch buffers.
  • Gather only referenced source batches, with null-safe nested gathering and the existing Arrow kernels for other encodings. Avoid redundant null processing for primitive payloads and preserve the no-null kernel when output needs no padding.
  • Avoid recoalescing already materialized multi-batch join output when its reliably known flat-column copy size reaches 2 MiB. Small outputs, compact builds, and fetch-clipped batches keep their existing buffering behavior; buffered prefixes and final build-row output preserve order. Nested/dictionary domains and view backing buffers do not inflate this copy-size test.
  • Support computed and composite keys, dictionary keys/payloads, nested payloads, ordinary mark/semi/anti/outer joins, and batched perfect hashing.
  • Opportunistically compact dictionary-heavy builds with otherwise fixed-width payloads when the copy can be admitted, amortizing dictionary-domain unification across repeated probes. Capacity, null-key safety, or admission failures leave the original batches and their reservation intact. Ineligible layouts still use the generic batched path; this is a representation policy, not a restriction on supported join types.
  • Preserve the current bounded final-build-row emission, dynamic-filter coordination, and reservation lifetime. Optional IN-list construction falls back to the hash-table predicate when it cannot be admitted or represented.
  • Add a 32-case physical-plan benchmark covering small builds, large payloads, computed/dictionary keys, dictionary/list payloads, tiny batches, shared slices, excess backing, low-match outer joins, and sustained probing of duplicate and unique dictionary domains, with perfect hashing enabled and disabled.

The batched path currently applies to ordinary hash joins above the 64 MiB compact-build threshold. Prepared reusable builds and null-aware joins retain their existing contiguous state contracts. This does not add spilling, change other join operators, or wire the representation into Comet.

Dictionary-heavy builds that successfully compact do not receive the batched representation's peak-memory benefit. Copy admission uses the existing payload-buffer estimate, not a complete bound on Arrow kernel scratch or process RSS; actual retained buffers are reconciled afterward.

Are these changes tested?

  • 1,510 physical join tests and 12 coalescer tests passed on 0e64452c64659be1c71e90bbec67974119a2ab8a, including computed-key, encoding, all-join-type, fetch, memory-limit, slice-retention, metadata-capacity, probe-preprocessing, coalescing, gather, dictionary-capacity, optional-compaction rollback, and wide-output order/copy regressions.
  • The canonical extended workspace suite passed: 12,317 Rust tests (8 ignored), all 525 SQL logic files, and 1,000 spill-pool fuzz iterations.
  • cargo clippy --workspace --all-targets --all-features -- -D warnings, formatting, license headers, and spelling checks passed.
  • Both exact-base and head release benchmark binaries passed all 32 synthetic row-count, build-ID sum, and reservation-release checks.

Local validation caveat: the dependency mirror does not yet serve the locked thiserror 2.0.21. Validation copies use thiserror and thiserror-impl at 2.0.20; the PR's Cargo.lock is unchanged. Both benchmark revisions use the same temporary dependency lockfile.

The concurrent SQL logic run also emitted a non-failing 79.4 MB allocator-versus-reservation drift diagnostic. The runner explicitly warns that file/consumer attribution is unreliable with concurrent files; this is not evidence attributing the drift to the join, nor a claim of complete allocator accounting.

Performance

Matched measurements against base 1d9be2e10794994fc1495c35363658dc7db0adf7 have driven fixes for tiny-batch and repeated dictionary workloads; correctness validation passes. The latest head's dictionary pilots passed, but the broader run exposed a material HJ Q11 slowdown and variation in compact-path controls. A subsequent fresh-process ABBA/BAAB diagnostic did not reproduce that slowdown, but unrelated Java workloads overlapped the diagnostic, so it does not settle performance attribution in either direction. Untimed plan captures confirmed that Q11 exercises batched ArrayMap output and Q20 is a compact ArrayMap control, with matching row and batch counts between revisions. The full run was also interrupted before its last TPC-H process when unrelated compilation started; all 22 persisted TPC-H value checks passed, but that run is incomplete performance evidence. All adverse, confounded, and interrupted runs are retained as investigation data. This remains a draft pending a quiet benchmark window, with no final speedup claim or performance table yet. Final evidence will compare the exact PR base with the final published head. The synthetic harness checks output row counts and matched build-ID sums, and separately records peak reservations, which are not process RSS.

Reproduction settings for the final comparison:

  • Shared x86-64 AMD EPYC-Milan host, 32 logical CPUs; rustc 1.98.1, release-nonlto, incremental compilation disabled, no extra Rust flags. No task-owned compilation or test jobs overlap the measurements. Host snapshots record unrelated activity; these are not dedicated-host timings.
  • Allocators are matched between revisions within each suite: the physical-plan Criterion harness uses the system allocator (glibc 2.39 on this host), and the SQL runner uses its existing default-feature mimalloc allocator.
  • The exact base receives only the identical benchmark source/registration and the matched dependency overlay described above; its production code is unchanged.
  • Synthetic: base/head/head/base, separate Criterion output directories, 1 s warmup, 2 s requested measurement, 20 samples per process. Slow cases may extend measurement time to retain all samples. Both perfect-hash settings are tested.
  • Existing HJ and TPC-H SQL suites: SF10 Parquet generated with tpchgen-cli 3.0.0, one part, ZSTD(1); 4 partitions, batch size 8,192, greedy memory pool limited to 16G. One warmup process per variant, then base/head/head/base with all three query iterations retained per process.
  • Timings aggregate each process's median, then the two process medians for each variant; no outlier removal. These are shared-host exploratory results. Reservation peaks are reported separately from timings and are not RSS measurements.

Build the runner with cargo build --profile release-nonlto -p datafusion-benchmarks --bin benchmark_runner. The synthetic benchmark command is:

cargo bench --profile release-nonlto -p datafusion-physical-plan \
  --bench hash_join_batches --features test_utils -- \
  --warm-up-time 1 --measurement-time 2 --sample-size 20 --noplot

Run each SQL suite with benchmark_runner hj or benchmark_runner tpch --format parquet, adding --scale-factor 10 --path DATA_PARENT --partitions 4 --batch-size 8192 --mem-pool-type greedy --memory-limit 16G --iterations 3 --result-mode none. DATA_PARENT contains tpch_sf10. Separate TPC-H warmups use --iterations 1 --result-mode persist to validate full result values before timing; the final results explain the Q11 tie-order and Q16 persistence-order qualifications.

Are there any user-facing changes?

No public API or configuration changes. Large ordinary hash joins may use less build memory by avoiding the full payload copy. Query results and the existing Arrow key-coercion requirements are unchanged.

Keep small builds compact, retain large builds with batch-local keys and bounded coalescing, and gather only referenced sources. Cover computed and composite keys, encoded and nested payloads, perfect hashing, memory admission, and bounded final output. Add a reproducible build-layout benchmark suite.
@github-actions github-actions Bot added the physical-plan Changes to the physical-plan crate label Sep 30, 2026
@codecov-commenter

codecov-commenter commented Sep 30, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 89.41450% with 273 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.61%. Comparing base (1d9be2e) to head (0e64452).

Files with missing lines Patch % Lines
...usion/physical-plan/src/joins/utils/multi_batch.rs 85.73% 58 Missing and 42 partials ⚠️
...ysical-plan/src/joins/hash_join/exec/build_data.rs 92.93% 27 Missing and 52 partials ⚠️
...fusion/physical-plan/src/joins/hash_join/stream.rs 84.97% 22 Missing and 10 partials ⚠️
...tafusion/physical-plan/src/joins/hash_join/exec.rs 89.73% 6 Missing and 17 partials ⚠️
datafusion/physical-plan/src/coalesce/mod.rs 83.62% 2 Missing and 17 partials ⚠️
datafusion/physical-plan/src/joins/array_map.rs 86.33% 16 Missing and 3 partials ⚠️
datafusion/physical-plan/src/joins/utils.rs 98.52% 0 Missing and 1 partial ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25889      +/-   ##
==========================================
+ Coverage   82.58%   82.61%   +0.03%     
==========================================
  Files        1144     1146       +2     
  Lines      441510   444014    +2504     
  Branches   441510   444014    +2504     
==========================================
+ Hits       364615   366843    +2228     
- Misses      54844    54980     +136     
- Partials    22051    22191     +140     

☔ 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.

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.

2 participants