Describe the bug
JoinSelection is not safe to run on a plan that has already been through FilterPushdown(Post): it re-decides the build side of every HashJoinExec, not only those still in PartitionMode::Auto, and when the decision flips it calls HashJoinExec::swap_inputs, which (correctly, since #23078) refuses once a dynamic filter has been constructed:
Internal error: Cannot swap HashJoinExec inputs after dynamic filters have been constructed.
Optimizer rules that reorder join inputs must run before FilterPushdown::new_post_optimization()
The guard lives in the operator; the rule has no corresponding check, so instead of leaving such a join alone it fails the whole plan.
To Reproduce
In join_selection.rs, statistical_join_selection_subrule handles all three modes:
PartitionMode::Auto → try_collect_left(hash_join, false, ..)
PartitionMode::CollectLeft → try_collect_left(hash_join, true, ..) (threshold ignored, may still swap)
PartitionMode::Partitioned → should_swap_join_order then swap_inputs
So a join that the first pass left as CollectLeft or Partitioned is re-evaluated on the next pass. The decision is statistics based, and after FilterPushdown(Post) the probe side carries a dynamic filter, so the inputs' statistics can differ from the first pass and the swap decision can flip.
Any downstream that re-runs the default physical optimizer on an already optimized plan hits this. We hit it in production planning (a LEFT JOIN of a large table against a small seed table, re-optimized after wrapping the plan in a writer sink); the failure is cost based, so unrelated edits to the query flip it between working and failing. It only reproduces with datafusion.optimizer.enable_join_dynamic_filter_pushdown = true (the default).
Unit level: build a HashJoinExec in PartitionMode::CollectLeft whose left input has larger statistics than its right, call set_dynamic_filter on it (as hash_join/exec.rs does in its own swap_inputs test), and run JoinSelection on it. try_collect_left(.., ignore_threshold = true, ..) decides to swap and returns the error above.
Expected behavior
JoinSelection should skip a HashJoinExec that already carries a dynamic filter (the join is past the point where its inputs may be reordered) and leave the plan unchanged for that node, instead of returning an error. A kept but suboptimal build side is the right trade here; carrying the dynamic filter across a swap would mean remapping its column references.
Additional context
The invariant is already stated in the swap_inputs doc comment and enforced there. This makes the rule honor it on its side, which also makes JoinSelection idempotent on its own output (related: #25688, #25585).
I can open the PR.
Describe the bug
JoinSelectionis not safe to run on a plan that has already been throughFilterPushdown(Post): it re-decides the build side of everyHashJoinExec, not only those still inPartitionMode::Auto, and when the decision flips it callsHashJoinExec::swap_inputs, which (correctly, since #23078) refuses once a dynamic filter has been constructed:The guard lives in the operator; the rule has no corresponding check, so instead of leaving such a join alone it fails the whole plan.
To Reproduce
In
join_selection.rs,statistical_join_selection_subrulehandles all three modes:PartitionMode::Auto→try_collect_left(hash_join, false, ..)PartitionMode::CollectLeft→try_collect_left(hash_join, true, ..)(threshold ignored, may still swap)PartitionMode::Partitioned→should_swap_join_orderthenswap_inputsSo a join that the first pass left as
CollectLeftorPartitionedis re-evaluated on the next pass. The decision is statistics based, and afterFilterPushdown(Post)the probe side carries a dynamic filter, so the inputs' statistics can differ from the first pass and the swap decision can flip.Any downstream that re-runs the default physical optimizer on an already optimized plan hits this. We hit it in production planning (a
LEFT JOINof a large table against a small seed table, re-optimized after wrapping the plan in a writer sink); the failure is cost based, so unrelated edits to the query flip it between working and failing. It only reproduces withdatafusion.optimizer.enable_join_dynamic_filter_pushdown = true(the default).Unit level: build a
HashJoinExecinPartitionMode::CollectLeftwhose left input has larger statistics than its right, callset_dynamic_filteron it (ashash_join/exec.rsdoes in its ownswap_inputstest), and runJoinSelectionon it.try_collect_left(.., ignore_threshold = true, ..)decides to swap and returns the error above.Expected behavior
JoinSelectionshould skip aHashJoinExecthat already carries a dynamic filter (the join is past the point where its inputs may be reordered) and leave the plan unchanged for that node, instead of returning an error. A kept but suboptimal build side is the right trade here; carrying the dynamic filter across a swap would mean remapping its column references.Additional context
The invariant is already stated in the
swap_inputsdoc comment and enforced there. This makes the rule honor it on its side, which also makesJoinSelectionidempotent on its own output (related: #25688, #25585).I can open the PR.