Skip to content

feat(pwmj): support left/right mark joins - #25467

Merged
jayzhan211 merged 2 commits into
apache:mainfrom
SubhamSinghal:pwmj-mark-joins
Sep 20, 2026
Merged

jayzhan211 merged 2 commits into
apache:mainfrom
SubhamSinghal:pwmj-mark-joins

Conversation

@SubhamSinghal

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Part of #17427.

Rationale for this change

PiecewiseMergeJoinExec supported every Semi/Anti existence join but rejected LeftMark/RightMark, falling back to NestedLoopJoinExec. Mark joins answer the same per-row existence question Semi/Anti already answer -- they just keep every row instead of filtering by it, attaching the answer as a boolean mark column. No algorithmic reason forces a single range predicate onto the slower path for these.

What changes are included in this PR?

  • LeftMark served by ExistencePWMJStream (existence_join.rs): same watermark as LeftSemi/LeftAnti, but the final pass emits every buffered row with a mark column built from it instead of slicing.
  • RightMark served by RightExistencePWMJStream (right_existence_join.rs): buffered-extreme comparison factored into a shared matched_mask; RightMark keeps every streamed row and appends the comparison as mark, folding NULL to false (mark is documented to never be NULL).
  • PiecewiseMergeJoinExec::try_new and the physical planner no longer reject/exclude Mark joins from the PWMJ path.
  • RightMark's maintains_input_order corrected: appending a column doesn't reorder rows, so it now matches RightSemi/RightAnti instead of a placeholder false.
  • Documented (in try_new) and pinned (SLT) an invariant: null-aware LeftMark (scalar NOT IN) can never reach PWMJ, since it always requires a non-empty equi-join on, while PWMJ only activates when on is empty. PiecewiseMergeJoinExec has no null_aware field and always builds a non-nullable mark, so this must hold.
PWMJ NestedLoopJoin speedup
LeftMark, all_match 246 µs 77.7 ms ~316×
LeftMark, no_match 224 µs 78.6 ms ~350×
LeftMark, half_match 245 µs 78.6 ms ~320×
RightMark, all_match 20.5 µs 78.8 ms ~3845×
RightMark, no_match 20.4 µs 78.9 ms ~3866×
RightMark, half_match 20.5 µs 78.6 ms ~3834×

What is the testing strategy for this PR?

  • Unit tests (existence_join.rs, right_existence_join.rs): join_left_mark, join_right_mark, all-NULL-buffered-side cases for both, RightMark NaN/-0.0 handling, a RightMark NULL-streamed-key case against a non-null buffered extreme, and a RightMark coalescer order-preservation case.
  • pwmj.slt: LeftMark has no SQL syntax of its own, so it's reached via EXISTS inside a disjunction (x > 100 OR EXISTS (...)), with EXPLAIN pinning PiecewiseMergeJoin: join_type=LeftMark. Covers real NULL buffered keys and List/LargeList/FixedSizeList/Dictionary keys. Also pins that a size-skewed range mark join never gets swapped to RightMark (PWMJ's swap_inputs is unimplemented) and that the null-aware NOT IN invariant above holds, plan and result both.
  • RightMark has no SQL path today (no optimizer rule constructs it), so it's covered only by unit tests and the fuzz test below.
  • fuzz_pwmj_matches_nested_loop (join_fuzz.rs, extended_tests): both Mark types added to the differential fuzz matrix against a NestedLoopJoinExec oracle, compared as (id, mark) pairs across randomized i32/f64 inputs, operators, partition counts, and NULL/duplicate keys.
  • Proto roundtrip tests updated: the wire-rejection test no longer expects Mark types to fail (no wire format change was needed -- the shared JoinType enum already covered them), and a new test executes both the original and roundtripped plan and compares output row-for-row, including mark.
  • cargo fmt, cargo clippy, and cargo doc all pass clean on the touched crates.

Are there any user-facing changes?

Yes: queries that decorrelate to a LeftMark/RightMark join with a single range predicate (no equi-join key) and no other filter now use PiecewiseMergeJoinExec instead of NestedLoopJoinExec -- a plan-shape and performance change, not a semantic one. No public API changes.

@SubhamSinghal

Copy link
Copy Markdown
Contributor Author

benchmark PR: #24955

@github-actions github-actions Bot added core Core DataFusion crate sqllogictest SQL Logic Tests (.slt) proto Related to proto crate physical-plan Changes to the physical-plan crate labels Sep 18, 2026
@codecov-commenter

codecov-commenter commented Sep 18, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 87.77778% with 33 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.37%. Comparing base (3b16a3d) to head (8885ca1).
⚠️ Report is 14 commits behind head on main.

Files with missing lines Patch % Lines
...joins/piecewise_merge_join/right_existence_join.rs 87.50% 5 Missing and 20 partials ⚠️
...ysical-plan/src/joins/piecewise_merge_join/exec.rs 89.74% 1 Missing and 3 partials ⚠️
...n/src/joins/piecewise_merge_join/existence_join.rs 86.66% 2 Missing and 2 partials ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25467      +/-   ##
==========================================
+ Coverage   82.34%   82.37%   +0.02%     
==========================================
  Files        1137     1138       +1     
  Lines      432713   433669     +956     
  Branches   432713   433669     +956     
==========================================
+ Hits       356338   357236     +898     
- Misses      54842    54862      +20     
- Partials    21533    21571      +38     

☔ 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

Copy link
Copy Markdown
Contributor Author

@comphead @kumarUjjawal PR for left and right mark join

@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 , some minor suggestions

Comment on lines +484 to +486
let mut mark = vec![false; len];
mark[min_marked..].fill(true);
Arc::new(BooleanArray::from(mark))

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.

Suggested change
let mut mark = vec![false; len];
mark[min_marked..].fill(true);
Arc::new(BooleanArray::from(mark))
let mut mark = BooleanBufferBuilder::new(len);
mark.append_n(min_marked, false);
mark.append_n(len - min_marked, true);
Arc::new(BooleanArray::new(mark.finish(), None))

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

/// column is documented never to be NULL; see `JoinType::LeftMark`).
fn mark_streamed_batch(&self, batch: &RecordBatch) -> Result<RecordBatch> {
let mark: ArrayRef = match self.matched_mask(batch)? {
None => Arc::new(BooleanArray::from(vec![false; batch.num_rows()])),

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.

Suggested change
None => Arc::new(BooleanArray::from(vec![false; batch.num_rows()])),
None => Arc::new(BooleanArray::new(
BooleanBuffer::new_unset(batch.num_rows()),
None,
)),

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

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

Review plus a simplification pass. No correctness bug found. I checked the risky plumbing rather than trusting it: EquivalenceGroup::join and updated_right_ordering_equivalence_class both leave right-side column indices unoffset for RightMark, so the corrected maintains_input_order = [false, true] is sound; build_join_schema declares mark nullable, so a null-free mark array is accepted; and boolean_mask_from_filter strips nulls, so RightMark cannot emit a NULL mark. Extracting matched_mask is a real dedup win, and dropping the LeftMark | RightMark => internal_err! arm makes execute's match exhaustive over all ten JoinType variants, so a new variant is now a compile error instead of a runtime one.

Findings inline. The substantive one is the try_new comment: its central claim is false, and it is the kind of comment a future maintainer will act on.

Most of the rest is test and comment volume. This PR adds 354 slt lines, 7 EXPLAIN pins (25 to 32 query TT in pwmj.slt), and a 75-line proto execution test for a feature whose semantics are one boolean column. Several of the new blocks say in their own comments that they re-cover ground already covered, or pin the absence of a rule that structurally cannot fire.

Not flagged, having checked: kv_float_{schema,batch,exec} correctly follows the established kv_<type>_* triple used by list/struct/dict; the new tables going undropped matches pwmj.slt's own convention (68 CREATE, 0 DROP); RightMark's Rust-only tests genuinely cannot move to slt, since no rule constructs RightMark.

One note on verification: cargo test 'extended_tests' was skipped on this run, so the new fuzz_pwmj_matches_nested_loop rows (the strongest evidence in the PR) have not actually executed in CI. I attempted to run them locally and ran out of disk. Worth confirming before merge.

Comment on lines +323 to +337
// There is no `null_aware` parameter here, unlike `HashJoinExec::try_new`: this
// constructor cannot express a null-aware mark join (see `JoinType::LeftMark`'s
// scalar-`NOT IN` variant, where `mark` is nullable), and `mark_streamed_batch`/
// `emit_matched` always build a non-nullable `mark` column. That is only sound
// because a null-aware mark join can never reach here: `decorrelate_predicate_subquery`
// only sets `null_aware` when the whole predicate is pure hash-equality with no
// residual (`mark_filter_is_hashable_only`), which means `join_on` is always
// non-empty for one -- and the PWMJ branch in `physical_planner.rs` only ever
// constructs this exec when `join_on` is empty. If either side of that ever changes
// (the PWMJ gate growing to accept a residual equijoin condition alongside a range
// predicate -- see the `TODO` on that branch -- or decorrelation producing
// `null_aware` from something other than a pure-equality predicate), this invariant
// breaks silently: a null-aware mark join would compute a plain boolean `mark` where
// SQL requires `NULL` (`UNKNOWN`), and nothing here would notice.
//

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 comment's central claim is false. It says that if the invariant breaks, "a null-aware mark join would compute a plain boolean mark where SQL requires NULL ... and nothing here would notice." The planner already notices, unconditionally:

// physical_planner.rs:1609-1614
if *null_aware && join_on.is_empty() {
    return plan_err!(
        "null_aware {join_type} join requires equi-join keys, but the join has none"
    );
}

That runs before the let join = if join_on.is_empty() { ... } branch that constructs this exec, so for null_aware == true either join_on is non-empty (and the PWMJ branch is not taken) or it is empty (and planning fails). PWMJ cannot see a null-aware join, by explicit guard rather than by emergent reasoning.

So the 16 lines reconstructing decorrelate_predicate_subquery's behavior are both unnecessary and misleading: they tell a maintainer the property is unguarded when it is enforced five statements earlier in the same function. Two lines pointing at that guard carry the whole argument, and stay true if decorrelation changes.

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

// collecting it. `RightMark` decides its `mark` column with the exact same
// comparison as `RightSemi`/`RightAnti` -- it just keeps every row instead of
// filtering by it -- so it takes the same path.
JoinType::RightSemi | JoinType::RightAnti | JoinType::RightMark => {

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 hardcodes the right-existence set while required_input_ordering, input_distribution_requirements, and benefits_from_input_partitioning all derive it from is_supported_right_existence_join. Four sites, three sharing a predicate and one spelling it out.

The two have to agree: this arm selects RightExistencePWMJStream, which reads the buffered side unordered and multi-partition, and the predicate is what asks for those relaxed requirements. If they drift, a join type gets routed to a stream whose input requirements were never requested, which is a wrong-answer bug rather than a compile error. Call the predicate here too.

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

Comment on lines 23 to 26
// `RightMark` belongs here too: deciding its mark column is the same one-key comparison as
// `RightSemi`/`RightAnti`, just kept instead of used to filter, so it needs no more of the
// buffered side than they do.
pub(super) fn is_supported_right_existence_join(join_type: JoinType) -> bool {

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 is_existence_join and is_supported_existence_join gone, "supported" no longer distinguishes anything: every existence join is supported now. The name reads as though some right existence join is still rejected. is_right_existence_join says what it tests.

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

Comment on lines +479 to +486
/// Builds the `LeftMark` `mark` column from the watermark: `false` for the unmatched prefix
/// `[0, min_marked)`, `true` for the matched suffix `[min_marked, len)` -- the same split
/// `LeftSemi`/`LeftAnti` slice the buffered batch on, just kept as one column instead of used
/// to drop rows.
fn mark_column(len: usize, min_marked: usize) -> ArrayRef {
let mut mark = vec![false; len];
mark[min_marked..].fill(true);
Arc::new(BooleanArray::from(mark))

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.

Once @jayzhan211's BooleanBufferBuilder form lands, this is four lines with one caller. Inline it into emit_matched and drop the helper plus its four-line doc comment, which restates the arm that calls it.

Supporting their suggestion rather than repeating it: BooleanBuffer::new_unset/BooleanBufferBuilder is already the house pattern for exactly this in the sibling mark implementation, including nested_loop_join.rs:3226 using BooleanBuffer::new_unset(right_batch.num_rows()) for the same "no buffered match, mark everything false" case that mark_streamed_batch builds with vec![false; n]. The current vec![false; len] plus fill(true) plus bit-packing is three passes and a byte-per-row temporary where the builder is one pass.

Comment on lines +428 to +431
// `LeftAnti`: the unmarked prefix, which includes every null-keyed row --
// nulls sort first and the watermark never drops below the buffered 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.

This is a three-way dispatch with two arms named and the third left as _, documented as LeftAnti in a comment. Spell it JoinType::LeftAnti and let a fourth join type reaching this stream be a compile error rather than silently taking the anti slice. Same shape as filter_streamed_batch's _ => for RightAnti, which is worth the same treatment.

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

Comment on lines +1066 to +1076
/// `roundtrip_test`/`roundtrip_test_and_return` only compare the `Debug` string of the
/// before/after plans -- which, per their own doc comment, "often isn't sufficient to
/// guarantee that no information is lost during serde because the string representation of
/// a plan often only shows a subset of state". `LeftMark`/`RightMark` add no new field to
/// encode (`join_type` already selects them from the shared proto enum, see
/// `join_type_to_proto`/`join_type_from_proto`), so the real risk is not a missing wire field
/// but a decoded plan that behaves differently at execution time. This actually executes both
/// the original and the roundtripped plan over real data and compares their output batches
/// row for row, including the `mark` column.
#[tokio::test]
async fn roundtrip_piecewise_merge_join_mark_executes_correctly() -> 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.

Drop this. roundtrip_piecewise_merge_join above already round-trips both Mark types, and its field-count assertion (4 for LeftMark, 3 for RightMark) is by itself proof that join_type survived the wire: no other join type produces those widths, and the exec derives everything else from the six constructor arguments inside try_new. Executing both Mark types over real data is what fuzz_pwmj_matches_nested_loop does, across randomized inputs, operators and partition counts, rather than one hand-built four-row case.

The doc comment argues the Debug-string comparison is too weak, which is true in general, but the PR description also establishes there is no new wire field here. 75 lines and three new imports (Int64Array, RecordBatch, MemorySourceConfig) for the one bit the field-count assertion already pins.

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

Comment on lines +553 to +563
# `PiecewiseMergeJoinExec::swap_inputs` is unimplemented (`todo!()`), and the physical
# optimizer's statistics-driven swap subrule (`join_selection.rs`) only ever downcasts to
# `HashJoinExec`/`CrossJoinExec`/`NestedLoopJoinExec` -- it does not consider
# `PiecewiseMergeJoinExec` at all. So unlike the equijoin case in `mark_join_matrix.slt`
# (where a size-skewed `LeftMark` on `HashJoinExec` gets swapped to `RightMark` with inputs
# flipped), a size-skewed range mark join through PWMJ has nothing to trigger that swap and
# must stay `LeftMark` regardless of which side is bigger. This pins that: if
# `swap_inputs` is ever implemented and wired in, this plan changing to `RightMark` is the
# signal to add real `RightMark` coverage here rather than an accidental behavior change.
statement ok
CREATE TABLE pwmj_mark_swap_l(k INT) AS VALUES (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 pins the absence of a rewrite that cannot fire. As the comment itself establishes, join_selection.rs only ever downcasts to HashJoinExec/CrossJoinExec/NestedLoopJoinExec, so no amount of size skew reaches PiecewiseMergeJoinExec::swap_inputs. The cost is a 1000-row generate_series table and a DataSourceExec: partitions=4, partition_sizes=[1, 0, 0, 0] pin that any change to target_partitions defaults or generate_series partitioning will break, for a signal that a todo!() nobody calls is still uncalled.

If the concern is that implementing swap_inputs later silently produces a RightMark with no coverage, the durable place to state that is a comment on swap_inputs itself.

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

Comment on lines +594 to +607
# `PiecewiseMergeJoinExec::try_new` has no `null_aware` parameter and always builds a
# non-nullable `mark` column, which is only sound because a null-aware `LeftMark` join (the
# one scalar `NOT IN` plans, per `JoinType::LeftMark`'s doc) can never reach the PWMJ branch:
# `decorrelate_predicate_subquery` only sets `null_aware` for a pure hash-equality predicate,
# which always has a non-empty `join_on`, while PWMJ only activates when `join_on` is empty.
# This pins the query that null-aware coverage in `null_aware_mark_join.slt` already uses
# (with `enable_piecewise_merge_join` on, unlike that file) and asserts it still plans to
# `HashJoinExec ... null_aware`, not `PiecewiseMergeJoin`: if that invariant is ever broken
# (the PWMJ gate loosened to accept a residual equijoin condition, or decorrelation start
# setting `null_aware` for something other than a pure-equality predicate), this is the
# query that should end up on PWMJ and lose its nullable `mark` -- so this plan changing is
# the signal, not a query that never had a chance to reach PWMJ in the first place.
statement ok
CREATE TABLE pwmj_null_aware_l(id INT) AS VALUES (1), (2), (NULL);

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.

Redundant with the planner guard (see my comment on exec.rs:323). physical_planner.rs:1609 turns a null-aware join with empty join_on into a plan_err!, so the outcome this block guards against cannot be reached, and the query pinned here is the one the comment admits null_aware_mark_join.slt already covers.

The one thing it adds over that file is running with enable_piecewise_merge_join = true (which defaults to false). That is worth a sentence, not two tables, an EXPLAIN and a duplicated three-valued result assertion: the invariant lives in the plan_err!, and deleting that guard already fails its own tests.

Comment on lines +1481 to +1505
# Existence joins: LeftMark with List/LargeList/FixedSizeList keys
# ------------------------------------------------------------------

statement ok
CREATE TABLE ex_list_l(id INT, v INT[]);

statement ok
INSERT INTO ex_list_l VALUES (1, [1]), (2, [2]), (3, [4]), (4, [5]);

statement ok
CREATE TABLE ex_list_r(v INT[]);

statement ok
INSERT INTO ex_list_r VALUES ([2]), ([4]);

# `>` : `EXISTS r WHERE l.v > r.v` holds iff `l.v` beats the smallest `r.v` ([2]), which only
# [4] (id 3) and [5] (id 4) do.
query I
SELECT l.id FROM ex_list_l l
WHERE l.id > 100 OR EXISTS (SELECT 1 FROM ex_list_r r WHERE l.v > r.v) ORDER BY 1;
----
3
4

# `LargeList`/`FixedSizeList` go through the same generic comparator as `List` above

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.

These roughly 100 lines cover key types through a code path that cannot be key-type sensitive, and the PR says so twice. pwmj.slt:1469: "LeftMark shares the exact same buffered-side comparator (JoinKeyComparator...) as LeftSemi/LeftAnti above -- only emit_matched's final step differs (slice vs. mark column)". That final step is buffered_batch.columns().to_vec() plus a boolean column, which cannot observe whether the key was List, LargeList, FixedSizeList or a dictionary. The comparator itself is already covered for all four types by the LeftSemi/LeftAnti cases above. Keep one nested-key LeftMark case as a smoke test and drop LargeList, FixedSizeList and the dictionary variant.

Same argument shrinks the EXPLAIN volume. pwmj.slt goes from 25 to 32 query TT, and five of the new ones pin the same shape (FilterExec ... OR mark@1 over ProjectionExec over PiecewiseMergeJoin: join_type=LeftMark) with only table and operator names differing: here, at 1505, and the multi-NULL (1224) and -0.0 (1308) blocks, whose comments likewise say they re-pin outcomes already covered for LeftSemi/LeftAnti. Keep the one canonical LeftMark EXPLAIN near the top of the section and make the rest result-only. Every unrelated planner change currently has to update all of them.

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

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

Thanks @SubhamSinghal no blockers for PR, please address the suggestions and we would be good to go

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

LGTM 🚀

@jayzhan211
jayzhan211 added this pull request to the merge queue Sep 20, 2026
@jayzhan211

Copy link
Copy Markdown
Contributor

Thanks @SubhamSinghal @kumarUjjawal @comphead 🚀

Merged via the queue into apache:main with commit d20936c Sep 20, 2026
41 checks passed
@SubhamSinghal
SubhamSinghal deleted the pwmj-mark-joins branch September 20, 2026 07:25
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 proto Related to proto crate sqllogictest SQL Logic Tests (.slt) v56.0.0

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants