zhuqi-lucas opened a new issue, #26106: URL: https://github.com/apache/datafusion/issues/26106
### 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. -- 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]
