Skip to content

perf(pwmj):settle never-matching streamed rows against the buffered extreme… - #25840

Open
SubhamSinghal wants to merge 6 commits into
apache:mainfrom
SubhamSinghal:classic-pwmj-extreme-key-binary-search
Open

SubhamSinghal wants to merge 6 commits into
apache:mainfrom
SubhamSinghal:classic-pwmj-extreme-key-binary-search

Conversation

@SubhamSinghal

@SubhamSinghal SubhamSinghal commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

It doesn't close #18221. On the issue's 100K × 100K shape, the join's time goes to its huge output, so this change leaves it about the same. See the benchmarks below.

Rationale for this change

The classic PiecewiseMergeJoin stream (INNER / LEFT / RIGHT / FULL) sorts each streamed batch, then walks the sorted buffered side row by row to find each streamed row's first match. Two costs are avoidable:

  • Rows that can't match anything still pay full price. They are sorted, gathered and walked through the buffered side like every other row. This covers keys below every buffered key for <, and NULL keys. When most streamed rows can't match, which is common for a selective range predicate, the operator does all this work and outputs nothing for them.
  • The first match is found by a linear walk. Every match set is a suffix [k, len) of the sorted buffered side, so k can be found with a binary search instead of stepping over every non-matching buffered row.

What changes are included in this PR?

All changes are in joins/piecewise_merge_join/, plus a new benchmark.

  • Rows that can't match are removed before the sort (classic_join.rs). Because every match set is a suffix, a streamed row matches anything at all only if it matches the last buffered key.
    • One vectorized comparison against that key, cached once when the buffered side is ready, finds those rows.
    • They are emitted as unmatched rows for Right/Full and dropped for Inner/Left. Only the remaining rows are sorted, gathered and scanned.
    • When no rows were removed, the path is unchanged except for that one comparison per batch. When some were, the sort indices are mapped back to positions in the batch, so it is still gathered only once.
    • An empty buffered side, or one whose keys are all NULL, short-circuits without the comparison.
  • The comparison must agree exactly with the scan, or it would change results.
    • Flat keys use one apply_cmp, which normalizes -0.0/+0.0 just as the scan's JoinKeyComparator does.
    • Nested keys are decided with the scan's own comparator. apply_cmp orders NULL elements inside a key ascending, while the comparator applies the sort options, descending for </<=, at every nesting level.
    • The scan's nested ordering for </<= gives different results from NestedLoopJoinExec when keys contain NULL elements. That is a pre-existing bug on main, which I'll file and fix separately. This PR keeps PWMJ's current results unchanged.
  • Binary search for each row's first match (utils.rs).
    • first_match(lo, hi, matches) replaces the linear walk in resolve_classic_join.
    • It searches [previous row's match, len). The streamed batch is sorted, so each row's first match is at or after the previous row's.
    • The existence stream already did an inline binary search. It now calls the same helper, with no change in behavior.
    • The operator dispatch moves into two small helpers, matches_on_equal and is_match, used by both streams.
  • Streamed NULL keys never reach the scan now, since the new pre-sort step removes them. The scan's streamed-NULL skip was dead code and is replaced by a debug_assert. The buffered-side NULL skip stays.

Benchmarks

The pwmj suite from #25839 (100K buffered keys, 2M streamed rows, five match regimes × INNER/LEFT/RIGHT/FULL). Reproduce with:

cargo run -p datafusion-benchmarks --release --bin benchmark_runner -- pwmj
regime main (median) Inner Left Right Full
no_match: no streamed row matches 47–50 ms −95.9% (25×) −95.9% (24×) −94.0% (17×) −93.7% (16×)
null_heavy: half NULL, rest no_match 45–50 ms −95.7% (23×) −95.7% (23×) −94.0% (17×) −93.7% (16×)
half_match: half no_match, half selective 371–375 ms −9.2% −8.7% −8.2% −8.4%
selective: every row matches the 1–4 smallest buffered keys 694–699 ms −2.9% −3.9% −3.6% −3.2%
all_match: every row matches every buffered key 162–163 ms +0.1% −0.6% −0.0% +0.3%

What is the testing strategy for this PR?

  • first_match_agrees_with_linear_scan (utils.rs): exhaustive over every range, answer and start position up to 40. It checks the result against a linear scan and bounds the number of comparisons by ⌈log2(hi − lo + 1)⌉.
  • matchable_rows_agrees_with_scan (classic_join.rs): the pre-sort check against a brute-force version of the scan's own definition of a match. It covers 12 key types, 4 operators, and a normal, all-NULL and empty buffered side.
    • The key types: Int32, Float64 with ±0.0/NaN/±inf, Utf8, Utf8View, Binary, Dictionary, Decimal128, Date32, Timestamp, Boolean, and two List cases with NULL elements.
    • With the nested-key branch removed, this test fails.
  • join_right_less_than_signed_zero_prefilter_agrees_with_scan: end to end with -0.0/+0.0 keys.
  • Existing tests:
    • all PWMJ unit tests;
    • the NLJ differential fuzz test fuzz_pwmj_matches_nested_loop (--features extended_tests), covering Inner/Left/Right/Full, NULL keys, small batches and multiple partitions;
    • the PWMJ and join sqllogictests.

Are there any user-facing changes?

No API or configuration changes.

@github-actions github-actions Bot added the physical-plan Changes to the physical-plan crate label Sep 28, 2026
@SubhamSinghal

Copy link
Copy Markdown
Contributor Author

Benchmark PR: #25839

@codecov-commenter

codecov-commenter commented Sep 28, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 91.80952% with 43 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.67%. Comparing base (991fd23) to head (2627267).
⚠️ Report is 85 commits behind head on main.

Files with missing lines Patch % Lines
...lan/src/joins/piecewise_merge_join/classic_join.rs 91.56% 14 Missing and 26 partials ⚠️
...sical-plan/src/joins/piecewise_merge_join/utils.rs 95.65% 2 Missing ⚠️
...n/src/joins/piecewise_merge_join/existence_join.rs 80.00% 0 Missing and 1 partial ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25840      +/-   ##
==========================================
+ Coverage   82.55%   82.67%   +0.11%     
==========================================
  Files        1141     1147       +6     
  Lines      440209   446769    +6560     
  Branches   440209   446769    +6560     
==========================================
+ Hits       363409   369349    +5940     
- Misses      54829    55006     +177     
- Partials    21971    22414     +443     

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

@SubhamSinghal

SubhamSinghal commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor Author

@comphead @jayzhan211 can you help in reviewing this PR.

@SubhamSinghal SubhamSinghal changed the title perf:settle never-matching streamed rows against the buffered extreme… perf(pwmj):settle never-matching streamed rows against the buffered extreme… Sep 30, 2026

@jayzhan211 jayzhan211 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.

Thanks @SubhamSinghal, LGTM!

buffer_idx = first_match(buffer_idx, buffered_len, |idx| {
is_match(cmp.compare(row_idx, idx), match_on_equal)
});
debug_assert!(buffer_idx < buffered_len);

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.

Maybe internal error here?

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.

+1. In release builds nothing fails here. With buffer_idx == buffered_len, count is 0 and an empty batch is pushed, so for Right/Full the streamed row is neither matched nor emitted as unmatched and is silently lost. The debug_assert_eq! on streamed NULLs a few lines up has the same shape. A NULL key that reached the scan compares Less than every non-NULL buffered key (nulls_first), so it would match all of them. Both checks are cheap enough for internal_err!.

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.

Addressed in a15bdda

@comphead comphead 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.

Checked locally against NestedLoopJoinExec with a differential fuzz over random 1 to 7 row, unsorted streamed batches (Inner/Left/Right/Full, all four operators, batch sizes 2, 3 and 8192, one or two streamed partitions). No mismatches on this head for i32, f64 with -0.0/NaN, Utf8, lists without NULL elements and the other flat key types in matchable_rows_agrees_with_scan. PWMJ output for list keys with NULL elements is identical to main on the same inputs, as the description says.

  • fuzz_pwmj_matches_nested_loop builds the streamed side as one-row batches (pwmj_parts_exec), so a batch is either fully matchable or not at all. The new mixed path (rows split off, the rest sorted and mapped back through positions) is only reached by the small hand-written cases. Chunking the streamed side into random multi-row batches there is a cheap guard for it.
  • The description mentions a new benchmark and cargo bench -p datafusion --bench pwmj_classic_sql, but neither is in this diff. The suite is in #25839 and runs through benchmark_runner -- pwmj. Please update the reproduction command and link it.
  • The classic join section of the PiecewiseMergeJoinExec docs in exec.rs still describes the linear walk. A line on the pre-sort check against the last buffered key and the binary search would keep it accurate.
  • Nits: null_padded_streamed_batch keeps the old new_stream_batch parameter name. The "settled before the sort, unmatched for Right/Full, dropped otherwise" explanation is repeated in fetch_stream_batch, split_off_never_matching, matchable_rows and the loop in resolve_classic_join, so a pointer is enough in three of them. Once the nested-key issue is filed, please reference it next to the nested branch in matchable_rows so the branch can go when the comparator is fixed.

// NULL keys never match and sort first, so the scan starts past the buffered ones.
// Streamed NULL keys never get here: `matchable_rows` settled them before the sort.
debug_assert_eq!(stream_values[0].null_count(), 0);
if !batch_process_state.processed_null_count {

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.

With streamed NULLs gone from the scan, processed_null_count only seeds buffer_idx with the buffered NULL count once per stream batch. Seeding start_buffer_idx right after batch_process_state.reset() in fetch_stream_batch removes the flag, its reset and this first-call branch. I tried it locally and the PWMJ unit tests and the fuzz still pass.

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.

Addressed in a15bdda

/// holds, so a check built on that would let the `+0.0` row through as a candidate
/// that the scan then rejects; with SQL semantics both rows are unmatched.
#[tokio::test]
async fn join_right_less_than_signed_zero_prefilter_agrees_with_scan() -> 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.

This overlaps the float64 case of matchable_rows_agrees_with_scan and the f64 pass of fuzz_pwmj_matches_nested_loop. A classic RIGHT/FULL case in piecewise_merge_join_matrix.slt Part 1 would replace it and be checked against NestedLoopJoin at batch sizes 1, 2, 100 and 8192, instead of a hard-coded snapshot. The matrix only has -0.0 for existence joins today. The same file could take one unsorted streamed batch mixing matching, non-matching and NULL keys to cover the positions remap. I tried both cases across the matrix locally and they pass on this head.

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.

Addressed in a15bdda

&sort_to_indices(buffered.as_ref(), Some(sort_options), None)?,
None,
)?;
let extreme = match sorted.len().checked_sub(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 copies the extreme computation from collect_buffered_side, so the test would not notice a bug in the production version. A small shared function, for example fn buffered_extreme(values: &ArrayRef) -> Result<Option<ColumnarValue>>, would remove the duplicate and make the test exercise the real code.

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.

Addressed in a15bdda

@github-actions github-actions Bot added core Core DataFusion crate sqllogictest SQL Logic Tests (.slt) labels Oct 1, 2026

@jayzhan211 jayzhan211 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.

There is one issue from the new commits

fn buffered_extreme(values: &ArrayRef) -> Result<Option<ColumnarValue>> {
Ok(match values.len().checked_sub(1) {
Some(last) if values.is_valid(last) => Some(ColumnarValue::Scalar(
ScalarValue::try_from_array(values, last)?,

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.

For Dictionary(_, Float*) keys the pre-filter disagrees with the scan on -0.0. apply_cmp normalizes the streamed array (dictionary values included), but normalize_float_zero_scalar skips ScalarValue::Dictionary, so the extreme keeps -0.0 under total order while JoinKeyComparator treats it as 0.0. This gives an internal error when the filter lets a row through, and wrong rows when it drops a real match. The base commit returns the right results.

Repro:

set datafusion.optimizer.enable_piecewise_merge_join = true;
CREATE TABLE dl AS SELECT * FROM (VALUES (1, arrow_cast(-0.0, 'Dictionary(Int32, Float64)')), (2, arrow_cast(0.5, 'Dictionary(Int32, Float64)'))) t(lid, lv);
CREATE TABLE dr AS SELECT * FROM (VALUES (1, arrow_cast(0.0, 'Dictionary(Int32, Float64)')), (2, arrow_cast(0.25, 'Dictionary(Int32, Float64)'))) t(rid, rv);
SELECT * FROM dl JOIN dr ON dl.lv < dr.rv;
-- Internal error: PiecewiseMergeJoin: streamed row 1 reached the scan without a match

CREATE TABLE dl2 AS SELECT * FROM (VALUES (1, arrow_cast(-1.0, 'Dictionary(Int32, Float64)')), (2, arrow_cast(-0.0, 'Dictionary(Int32, Float64)'))) t(lid, lv);
CREATE TABLE dr2 AS SELECT * FROM (VALUES (1, arrow_cast(0.0, 'Dictionary(Int32, Float64)'))) t(rid, rv);
SELECT * FROM dl2 RIGHT JOIN dr2 ON dl2.lv >= dr2.rv;
-- got `NULL NULL 1 0`, expected `2 0 1 0` (NLJ and the base commit agree)

Fix (normalize the extreme before taking it as a scalar):

+use datafusion_common::utils::normalize_float_zero;
 fn buffered_extreme(values: &ArrayRef) -> Result<Option<ColumnarValue>> {
     Ok(match values.len().checked_sub(1) {
-        Some(last) if values.is_valid(last) => Some(ColumnarValue::Scalar(
-            ScalarValue::try_from_array(values, last)?,
-        )),
+        // `apply_cmp` normalizes `-0.0` only in flat float scalars, not inside a
+        // `ScalarValue::Dictionary`, so normalize the key before taking it.
+        Some(last) if values.is_valid(last) => {
+            let extreme = normalize_float_zero(&values.slice(last, 1));
+            Some(ColumnarValue::Scalar(ScalarValue::try_from_array(
+                &extreme, 0,
+            )?))
+        }
         _ => None,
     })
 }

Cases for matchable_rows_agrees_with_scan. They fail on the current head (dictionary_float64_neg_zero_min/full <: streamed row 0) and pass with the fix:

(
    // `-0.0` as the buffered extreme inside a dictionary: the smallest
    // key for `<`/`<=`, the largest for `>`/`>=`.
    "dictionary_float64_neg_zero_min",
    Arc::new(DictionaryArray::<Int32Type>::new(
        Int32Array::from(vec![0, 1]),
        Arc::new(Float64Array::from(vec![-0.0, 0.5])),
    )),
    Arc::new(DictionaryArray::<Int32Type>::new(
        Int32Array::from(vec![0, 1, 2]),
        Arc::new(Float64Array::from(vec![0.0, -0.0, 0.25])),
    )),
),
(
    "dictionary_float64_neg_zero_max",
    Arc::new(DictionaryArray::<Int32Type>::new(
        Int32Array::from(vec![0, 1]),
        Arc::new(Float64Array::from(vec![-1.0, -0.0])),
    )),
    Arc::new(DictionaryArray::<Int32Type>::new(
        Int32Array::from(vec![0, 1, 2]),
        Arc::new(Float64Array::from(vec![0.0, -0.0, -0.5])),
    )),
),

The root cause, normalize_float_zero_scalar not recursing into ScalarValue::Dictionary, is already on main: SELECT arrow_cast(0.0, 'Dictionary(Int32, Float64)') > arrow_cast(-0.0, 'Dictionary(Int32, Float64)') returns true. Worth filing on its own; fixing it there would also cover this.

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.

Thanks @jayzhan211. Addressed in 036502f

@jayzhan211 jayzhan211 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.

2 suggestions

)?;
let last = buffered_values.len() - 1;
let matchable = BooleanBuffer::collect_bool(num_rows, |row| {
stream_values.is_valid(row)

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.

is_valid only checks the physical null buffer, and a RunArray has none. So a run-end encoded NULL key passes as matchable, gets past the null_count() > 0 guard at :492, and the scan joins it with every buffered row (nulls_first → Less). main gives the same rows, so this is fine as a follow-up. But join_right_run_end_encoded_keys_empty_sliced_batch already feeds such a NULL (a2 = 0) and asserts nothing.

         let last = buffered_values.len() - 1;
+        // A run-end encoded NULL has no physical null buffer.
+        let nulls = stream_values.logical_nulls();
         let matchable = BooleanBuffer::collect_bool(num_rows, |row| {
-            stream_values.is_valid(row)
+            nulls.as_ref().is_none_or(|n| n.is_valid(row))
                 && is_match(cmp.compare(row, last), match_on_equal)
         });
-    if stream_values[0].null_count() > 0 {
+    if stream_values[0].logical_null_count() > 0 {

In the test, after importing batches_to_sort_string. This fails on this head and passes with the fix:

let (_, batches, _) =
    join_collect(left, right, on, Operator::Lt, JoinType::Right).await?;
// The NULL key (a2 = 0) and 0 (a2 = 3) match nothing.
assert_snapshot!(batches_to_sort_string(&batches), @r"
+----+----+----+----+
| a1 | b1 | a2 | b2 |
+----+----+----+----+
|    |    | 0  |    |
|    |    | 3  | 0  |
| 1  | 5  | 4  | 6  |
| 2  | 5  | 4  | 6  |
| 3  | 1  | 1  | 2  |
| 3  | 1  | 2  | 2  |
| 3  | 1  | 4  | 6  |
| 4  | 1  | 1  | 2  |
| 4  | 1  | 2  | 2  |
| 4  | 1  | 4  | 6  |
+----+----+----+----+
");

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.

Addressed in 2627267

&[sort_options],
NullEquality::NullEqualsNothing,
)?;
for row in 0..streamed.len() {

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 PR also fixes dictionary keys whose values are NULL (a valid key pointing at a NULL value). apply_cmp sees the logical NULL and sets the row aside, while main joins it with every buffered row. Nothing pins this, and the oracle here uses physical is_valid, so it would call such a row matchable. Fine as a follow-up; this passes on the current head:

+                    let nulls = streamed.logical_nulls();
                     for row in 0..streamed.len() {
-                        let expected = streamed.is_valid(row)
+                        let expected = nulls.as_ref().is_none_or(|n| n.is_valid(row))
(
    // A valid key pointing at a NULL value: logically NULL, physically valid.
    "dictionary_null_values",
    Arc::new(DictionaryArray::<Int32Type>::new(
        Int32Array::from(vec![0, 1, 2]),
        Arc::new(Int32Array::from(vec![Some(5), None, Some(1)])),
    )),
    Arc::new(DictionaryArray::<Int32Type>::new(
        Int32Array::from(vec![0, 1, 2, 3]),
        Arc::new(Int32Array::from(vec![Some(2), None, Some(6), Some(0)])),
    )),
),

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.

Addressed in 2627267

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

core Core DataFusion crate physical-plan Changes to the physical-plan crate sqllogictest SQL Logic Tests (.slt)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Check PiecewiseMergeJoin performance for large tables

4 participants