Repository navigation
fix: skip HashJoinExec nodes carrying dynamic filters in JoinSelection - #26123
siddubakka wants to merge 8 commits into
Conversation
apache#26106) ## Which issue does this PR close? - Closes apache#26106. ## Rationale for this change `JoinSelection` is not safe to run on plans that already carry dynamic filters, which occurs during re-optimization passes (e.g. downstream pipelines that wrap an already-optimized plan into a writer sink and re-run physical optimization). Once `FilterPushdown` has wired up a dynamic filter between the build side and probe side of a `HashJoinExec`, swapping the inputs via `swap_inputs` panics with: `Internal error: Cannot swap HashJoinExec inputs after dynamic filter has been constructed` The join's build side is already committed at that stage, and swapping inputs would invalidate the dynamic filter expressions that reference probe-side columns. `JoinSelection` should detect this and leave the `HashJoinExec` unchanged instead of failing the query. ## What changes are included in this PR? - In `JoinSelection::statistical_join_selection_subrule`, check if `!hash_join.dynamic_expressions_produced().is_empty()` and return `None` (leaving the plan unchanged). - In `can_swap_hash_join`, guard against swapping when dynamic expressions are produced. - In `hash_join_swap_subrule`, guard against swapping unbounded left inputs when dynamic expressions are produced. - Added tests in `join_selection.rs` verifying that `JoinSelection` skips `HashJoinExec` carrying dynamic filters in both `CollectLeft` and `Partitioned` modes. ## What is the testing strategy for this PR? Added `test_join_selection_skips_hash_join_with_dynamic_filter` in `datafusion/core/tests/physical_optimizer/join_selection.rs` verifying that `JoinSelection` leaves the plan unchanged for both `CollectLeft` and `Partitioned` modes without error. ## Are there any user-facing changes? No API changes. Fixes an internal error when re-optimizing plans that contain dynamic filters.
|
@siddubakka The CI has some fails, need to make it pass first. |
|
@zhuqi-lucas fixed the fmt and test failure, CI should be good now |
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #26123 +/- ##
==========================================
- Coverage 82.74% 82.74% -0.01%
==========================================
Files 1147 1147
Lines 449767 449771 +4
Branches 449767 449771 +4
==========================================
Hits 372160 372160
- Misses 54938 54941 +3
- Partials 22669 22670 +1 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
| lit(true), | ||
| )); | ||
|
|
||
| #[allow(deprecated)] |
There was a problem hiding this comment.
This attribute makes the current clippy --all-targets --workspace ... -D warnings job fail: Clippy 1.99 denies clippy::allow_attributes and explicitly suggests #[expect]. Replacing it with #[expect(deprecated)] preserves the intended scoped deprecation handling and also warns if the deprecation is later removed.
There was a problem hiding this comment.
@mikamikasuki thank you for your help i have made the changes now
|
@zhuqi-lucas fixed the clippy failure and the codecov issue (removed an unreachable branch that was dropping coverage). also added a test for it, so everything should be passing now. |
|
hey @zhuqi-lucas @xudong963, i pushed the latest updates yesterday. could one of you approve the workflow runs when you have a moment so the CI checks can kick off? thanks! |
|
|
| Ok(()) | ||
| } | ||
|
|
||
| #[rstest] |
There was a problem hiding this comment.
Could we add a regression that lets FilterPushdown::new_post_optimization() create the dynamic filter and probe-side consumer, then reruns JoinSelection and executes the plan? The current tests correctly exercise both guards, but manually calling with_dynamic_filter_expr does not verify the producer/consumer connection or query results.
- run cargo fmt to fix test_join_selection_skips_unbounded_hash_join_with_dynamic_filter signature - add end-to-end regression that lets FilterPushdown::new_post_optimization create the dynamic filter, then verifies JoinSelection leaves it unchanged and the plan executes with correct results
|
@rgbuilds hey thanks for the review, fixed the fmt failure on that test signature and yeah you're right about the manual with_dynamic_filter_expr test — added one that runs FilterPushdown::new_post_optimization to wire up a real filter, then checks JoinSelection leaves it alone and the plan still returns the right rows |
Thanks—the formatting issue is fixed, and the real One issue remains: both Could we give the left/build side larger reported statistics than the right/probe side, or otherwise force It would also be helpful to assert that the original left and right children are preserved. Then the test should fail without the guard and pass with it. |
|
@rgbuilds ah yeah missed that, both sides were 123 bytes so the swap never triggered. forced the build side bigger with a registry now and added ptr_eq checks on both children, should actually exercise the guard now |
Which issue does this PR close?
Rationale for this change
JoinSelectionis not safe to run on plans that already carry dynamic filters. This happens during re-optimization passes (e.g. downstream pipelines that wrap an already-optimized plan into a writer sink and re-run physical optimization).Once
FilterPushdownhas wired up a dynamic filter between the build side and probe side of aHashJoinExec, swapping the inputs viaswap_inputspanics with:The join's build side is already committed at that stage, and swapping inputs would invalidate the dynamic filter expressions that reference probe-side columns.
JoinSelectionshould detect this and leave theHashJoinExecunchanged instead of failing the plan.What changes are included in this PR?
JoinSelection::statistical_join_selection_subrule, check if!hash_join.dynamic_expressions_produced().is_empty()and returnNone(leaving the plan unchanged).can_swap_hash_join, guard against swapping when dynamic expressions are produced.hash_join_swap_subrule, guard against swapping unbounded left inputs when dynamic expressions are produced.datafusion/core/tests/physical_optimizer/join_selection.rsverifying thatJoinSelectionskipsHashJoinExeccarrying dynamic filters in bothCollectLeftandPartitionedmodes.What is the testing strategy for this PR?
Added
test_join_selection_skips_hash_join_with_dynamic_filterindatafusion/core/tests/physical_optimizer/join_selection.rsverifying thatJoinSelectionleaves the plan unchanged for bothCollectLeftandPartitionedmodes without error.Are there any user-facing changes?
No API changes. Fixes an internal error when re-optimizing plans that contain dynamic filters.