Skip to content

JoinSelection should skip HashJoinExec nodes that already carry a dynamic filter instead of failing in swap_inputs #26106

Description

@zhuqi-lucas

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.

Activity

  1. self-assigned this
    on Oct 7, 2026
  2. siddubakka commented on Oct 7, 2026

    @siddubakka

    take

  3. siddubakka commented on Oct 7, 2026

    @siddubakka

    take

  4. siddubakka commented on Oct 7, 2026

    @siddubakka

    @zhuqi-lucas I will raise a pr soon please feel free to review my work

  5. siddubakka commented on Oct 7, 2026

    @siddubakka

    @zhuqi-lucas I have opened PR #26123 to fix this. Please feel free to take a look!

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Labels

bugSomething isn't working

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions