SubhamSinghal opened a new pull request, #25840:
URL: https://github.com/apache/datafusion/pull/25840

   ## Which issue does this PR close?
   
   - Part of #18221 (and the PWMJ epic #17427).
   
   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
   
   `cargo bench -p datafusion --bench pwmj_classic_sql`, Apple M-series, tp=1. 
The baseline is `main` (`bf01e288e`): the same build with only 
`classic_join.rs` swapped back to `main`'s version, saved with `--save-baseline 
main`. Every change listed has p = 0.00.
   
   | 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.


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