Repository navigation
fix: sort-merge join filter columns follow the filter's own column order - #25489
Conversation
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25489 +/- ##
==========================================
+ Coverage 82.37% 82.42% +0.04%
==========================================
Files 1138 1138
Lines 433506 435493 +1987
Branches 433506 435493 +1987
==========================================
+ Hits 357102 358935 +1833
+ Misses 54850 54848 -2
- Partials 21554 21710 +156 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
There was a problem hiding this comment.
Thanks for splitting this out of #25217 — much easier to review on its own.
The diagnosis and the fix are correct. JoinFilter::swap negates each ColumnIndex.side in place and reuses the same schema Arc, so positions stay put while sides flip, and the intermediate schema stays positionally aligned with column_indices. The old two-pass collect rebuilt a left-first vector and zipped it against a schema that was no longer left-first. Walking column_indices once is the right fix, and it's a strict generalization: for a canonical left-first filter the new code yields the identical vector, so no existing plan can change behavior.
I also checked the second call site, which isn't visible from the diff and is load-bearing. In freeze_streamed_matched the locals left_columns / right_columns are streamed / buffered, and for JoinType::Right the probe side is the join's right input (exec.rs:238, probe_side), so the call site passes them swapped (materializing_stream.rs:1551-1553). That makes get_filter_columns's parameters join-side-relative in both branches, and I traced the Right case explicitly: it is broken on main in exactly the same way and is fixed by this change. Worth stating because nothing in the diff shows it.
On the bug class: I swept every consumer of filter.column_indices() under physical-plan/src/joins/ — bitwise_stream.rs:1239, nested_loop_join.rs:3433, utils.rs:1329, stream_join_utils.rs:276 (which even carries a comment asserting the ordering invariant), plus proto.rs:114 for serialization. All walk in order. asof_join.rs's column_indices is the join's output schema from build_join_schema, not a filter's, and it only tests for presence of a right-side column. So sort_merge_join/filter.rs was the last instance, and this closes the class rather than patching one case.
Two requests before merge, plus smaller notes inline.
1. The JoinSide::None arm doesn't match the code it cites.
The description says this "matches what the semi/anti/mark stream already does in bitwise_stream.rs" and that None columns "are skipped, as before". The ordering half matches; the None half doesn't — bitwise_stream.rs:1245 and :1253 return internal_err!("Unexpected JoinSide::None in filter"). Details inline; short version is that None in a filter's column_indices is an invariant violation, not a legitimate case, and skipping it defers the failure into an opaque RecordBatch::try_new arity error. Either match bitwise_stream.rs, or keep the skip and drop the "as before" framing.
2. The filter path through swap_inputs() has no test.
swap_inputs_swaps_the_projection (tests.rs:6735) already does exactly the right thing structurally — build the join, swap_inputs(), execute both, assert equal results — but it passes filter: None, so it covers projection renumbering only and never touches a filter. Since swap_inputs() is the one realistic way to reach this bug, extending that pattern with a real JoinFilter would cover the trigger end-to-end. That differential shape is also stronger than a snapshot: it asserts directly that swapping preserves semantics, which is the actual contract, instead of pinning one expected table.
On the Rationale — I think it undersells the fix.
It's honest that SQL-planned queries aren't affected, but it stops there, which left me deriving why the fix matters. The evidence is better than the description suggests, and it's first-hand rather than speculative:
- #13984, which added this method, gives the reason: to "enable implementation of (and experimentation with) custom join rules by downstream projects". It also noted up front that it shipped with no explicit unit tests.
- #17373 documented
HashJoinExec::swap_inputswith "This function is public so other downstream projects can use it to constructHashJoinExecwith right side as the build side" — and SMJ's own doc comment points at that text.
So this isn't unreachable dead code; it's a deliberately public extension point for downstream engines, which currently returns silently wrong rows whenever a filter is present and the misplaced columns share a type. For a library that's worse than a panic: downstream follows the documented interface, gets wrong data, and gets no signal. Saying that in the Rationale would make the PR much easier to evaluate. I confirmed in-tree reachability myself: no production caller, and it's an inherent method (not on any trait), so it can't be reached through dyn ExecutionPlan either.
Nit on the description's before/after table. The "actual" side shows two rows and omits | 30 | 6 | | | |. The buggy output has three rows too — b1 = 6 has no match on the right either way, so it's null-joined in both. Since that table is the description's main evidence it's worth having exact.
Adjacent, pre-existing, not this PR's job — flagging only because I was in the area. The !filter_columns.is_empty() guard at materializing_stream.rs:1563 uses "no filter columns" as a proxy for "no filter", so a filter with zero column references (1 = 1, say) would be skipped rather than evaluated. Identical on both sides of this diff.
| .filter_map(|col_index| match col_index.side { | ||
| JoinSide::Left => Some(Arc::clone(&left_columns[col_index.index])), | ||
| JoinSide::Right => Some(Arc::clone(&right_columns[col_index.index])), | ||
| JoinSide::None => None, |
There was a problem hiding this comment.
This is where the change diverges from bitwise_stream.rs, which the description cites as the model: bitwise_stream.rs:1245 and :1253 return internal_err!("Unexpected JoinSide::None in filter") rather than skipping.
I think erroring is the better behavior, because None in a filter's column_indices looks like an invariant violation rather than a legitimate case:
JoinSide::Noneexists for the mark column in a join's outputcolumn_indices(utils.rs:316,:328inbuild_join_schema), not a filter's.- Every construction site of a filter's
column_indicestreats it as impossible:projection_pushdown.rs:209,:217,:263,:642are allunreachable!("Mark join not supported");join_selection.rs:452andsort_pushdown.rs:888likewise. - Decisively: the intermediate schema is built from
column_indices, one field per entry (projection_pushdown.rs:211-219). ANoneentry can't even acquire a schema field.
So if one ever appeared, skipping yields a filter_columns shorter than f.schema(), and the failure surfaces a few lines later at materializing_stream.rs:1566 as an opaque column-count mismatch instead of a named internal error pointing at the real cause.
This is also the one line Codecov flags as uncovered, which is consistent with it being unreachable by construction.
I realize matching bitwise_stream.rs means returning Result<Vec<ArrayRef>>, which ripples to the two call sites, and there's a fair argument that's beyond a minimal fix. If you'd rather keep the skip, could you drop the "matches bitwise_stream.rs" / "as before" framing from the description and leave a comment here on why skipping is safe? As written the description reads as though the two implementations agree.
| // The filter's intermediate schema lists its columns in `column_indices` | ||
| // order, which need not put every left column before every right one | ||
| // (a filter swapped along with the join's inputs does the opposite). | ||
| f.column_indices() |
There was a problem hiding this comment.
Since you're rewriting this function anyway: could the doc comment say the parameters are join-side-relative, not streamed/buffered?
The caller has to re-map them. materializing_stream.rs:1551-1553 passes them swapped for JoinType::Right, because there the probe side is the join's right input (exec.rs:238) while the locals are named left_columns / right_columns after streamed/buffered. The correctness of this fix depends on that swap, but nothing here says so. A future caller inside the SMJ — where the locals carry exactly those names — would pass them positionally and be silently wrong for Right only, which is the same failure mode this PR is fixing.
Something like "callers pass arrays in join-side order, not streamed/buffered order; see the JoinType::Right call site in materializing_stream.rs" would earn its keep.
Separately, on the comment just above: it says what the order is, but not why the old code was wrong. The load-bearing fact — that the result is zipped against the filter's intermediate schema positionally, so a left-first vector against a right-first schema silently puts the wrong array in each slot — is in the PR description but not in the code. Having it here means the next reader doesn't need to find this PR to reconstruct the invariant.
| ])), | ||
| ); | ||
|
|
||
| let (_, batches) = join_collect_with_filter(left, right, on, filter, Left).await?; |
There was a problem hiding this comment.
Two suggestions on coverage.
A JoinType::Right variant. Right is the one join type in the materializing stream whose probe side is the join's right input, so it's where the join-side and streamed/buffered coordinate systems meet, and it takes the swapped-argument path at materializing_stream.rs:1551. I traced it: it's broken on main the same way and fixed by this change, but nothing pins it — so a future refactor of that if/else could reintroduce the bug on Right alone with the suite green. The existing filter tests (join_right_different_columns_count_with_filter, join_left_different_columns_count_with_filter, the mark/semi/anti ones) all use canonical left-first layouts, so Right x right-before-left is the highest-risk untested combination right now.
A test through swap_inputs(). This is the more valuable one, and there's already a template for it: swap_inputs_swaps_the_projection (tests.rs:6735) builds a join, calls swap_inputs(), executes both, and asserts the results are equal. It passes filter: None, so the filter path is untouched. Passing a real JoinFilter through that same shape would exercise the documented trigger end-to-end, and the differential assertion (swap must not change the result) is stronger than a snapshot because it states the contract rather than one expected table. It would also guard the rest of that path, which nothing does today.
|
Thanks @viirya , address them all! |
kosiew
left a comment
There was a problem hiding this comment.
Thanks for the fix. The change looks good overall, especially the explicit preservation of JoinFilter::column_indices order and the added coverage around right joins and swap_inputs().
I just have one small testing suggestion below.
| async fn join_left_with_filter_columns_right_before_left() -> Result<()> { | ||
| // select * | ||
| // from t2 | ||
| // left join t1 on t2.b1 = t1.b1 and t2.a2 > t1.a1 |
There was a problem hiding this comment.
Could we also add one three-column interleaving case, for example Left, Right, Left? The new tests cover the swapped Right, Left layout nicely, but this helper is meant to support arbitrary column_indices ordering. An interleaving case would help protect that broader contract from accidentally being changed back to grouping columns by side in the future.
…der (apache#25489) ## Which issue does this PR close? - No separate issue. Split out of apache#25217 so it can be reviewed on its own. ## Rationale for this change A `SortMergeJoinExec` with a join filter returns wrong rows when the filter's `column_indices` do not list every left column before every right column. `JoinFilter::swap` produces exactly that layout, so any plan that goes through the public `SortMergeJoinExec::swap_inputs()` is affected, as is any `JoinFilter` built by hand or by a custom optimizer rule. For example, a Left join on `t2.b1 = t1.b1 AND t2.a2 > t1.a1` whose filter indices are `[Right(a1), Left(a2)]` is evaluated as `a1 > a2`: expected actual | 10 | 4 | 1 | 4 | 7 | | 10 | 4 | | | | | 20 | 5 | | | | | 20 | 5 | 21 | 5 | 8 | When the misplaced columns share a type there is no error, only wrong results. Queries planned from SQL are not affected today: the physical planner always builds the filter with left columns first, and `JoinSelection` does not swap sort-merge joins. ## What changes are included in this PR? `get_filter_columns` (used by the materializing SMJ stream) collected all left columns, then all right columns, and the result was zipped against the filter's intermediate schema, which is in `column_indices` order. It now walks `column_indices` once and takes each column from the side it names. This matches what the semi/anti/mark stream already does in `bitwise_stream.rs`. Columns with `JoinSide::None` are skipped, as before. ## What is the testing strategy for this PR? New unit test `join_left_with_filter_columns_right_before_left` in `sort_merge_join/tests.rs`: a Left join whose filter lists a right column before a left one. It fails on `main` with the wrong rows shown above and passes with the fix. The existing `sort_merge_join` tests pass unchanged. ## Are there any user-facing changes? No API changes. Sort-merge joins with a filter in right-before-left column order now return correct results.
Which issue does this PR close?
Rationale for this change
A
SortMergeJoinExecwith a join filter returns wrong rows when the filter'scolumn_indicesdo not list every left column before every right column.JoinFilter::swapproduces exactly that layout, so any plan that goes throughthe public
SortMergeJoinExec::swap_inputs()is affected, as is anyJoinFilterbuilt by hand or by a custom optimizer rule. For example, a Leftjoin on
t2.b1 = t1.b1 AND t2.a2 > t1.a1whose filter indices are[Right(a1), Left(a2)]is evaluated asa1 > a2:When the misplaced columns share a type there is no error, only wrong results.
Queries planned from SQL are not affected today: the physical planner always
builds the filter with left columns first, and
JoinSelectiondoes not swapsort-merge joins.
What changes are included in this PR?
get_filter_columns(used by the materializing SMJ stream) collected all leftcolumns, then all right columns, and the result was zipped against the filter's
intermediate schema, which is in
column_indicesorder. It now walkscolumn_indicesonce and takes each column from the side it names.This matches what the semi/anti/mark stream already does in
bitwise_stream.rs. Columns withJoinSide::Noneare skipped, as before.What is the testing strategy for this PR?
New unit test
join_left_with_filter_columns_right_before_leftinsort_merge_join/tests.rs: a Left join whose filter lists a right columnbefore a left one. It fails on
mainwith the wrong rows shown above andpasses with the fix. The existing
sort_merge_jointests pass unchanged.Are there any user-facing changes?
No API changes. Sort-merge joins with a filter in right-before-left column
order now return correct results.