Skip to content

Use ProbeCompletion in NestedLoopJoinExec #26148

Description

@jayzhan211

Is your feature request related to a problem or challenge?

When a join shares one build side across probe partitions, the last partition to finish emits the build-side rows, for example the unmatched rows of a LEFT or FULL join. The rule is simple, but it is only correct if four things hold:

  1. Each partition reports once. In fix: NestedLoopJoinExec emits spurious unmatched-left rows with multiple probe partitions #22791, NLJ re-entered its emit state and decremented the counter a second time. The counter reached zero while other partitions were still probing, so a partition emitted unmatched left rows too early, and some rows appeared twice.
  2. The count covers every partition that shares the matches. In fix: refuse memory-limited NestedLoopJoin fallback for left-emission joins with a multi-partition probe side #24675 (Correctness guard: disable memory-limited NestedLoopJoin fallback for LEFT-family joins with multi-partition right side #22641), the spill path built a per-chunk JoinLeftData with a counter of 1. Every partition then thought it was last and emitted from its own partial bitmap. That returned wrong rows, for example 49 instead of 19 for a LEFT join.
  3. A partition that goes away still reports. In fix: cancel the NestedLoopJoin coordinated fallback when a partition is dropped unfinished #25004, a dropped NLJ partition never reported. No partition became last, and the others waited forever.
  4. The last partition reads the shared state only after its own decrement, with an ordering that makes the other partitions' writes visible. In fix: null-aware anti join could emit rows after a sibling probe partition saw NULL #25076, hash join read the shared "saw NULL" flag before decrementing a Relaxed counter, so NOT IN returned rows it must suppress.

#25076 introduced ProbeCompletion (joins/hash_join/probe_completion.rs) to make point 4 part of the type instead of a convention:

  • The shared facts are only available through the return value of report_completed(), so they cannot be read before the decrement.
  • The decrement is AcqRel.
  • The orderings are written once and checked with loom.

NestedLoopJoinExec still provides all four guarantees by hand, in two places with different orderings:

Moving NLJ onto ProbeCompletion brings:

Describe the solution you'd like

Use ProbeCompletion in NestedLoopJoinExec too:

  1. Move hash_join/probe_completion.rs to joins/probe_completion.rs.
  2. Add a call for a partition that stops early, for example abandon(). It still counts down, so nothing waits forever (point 3), and it tells the last partition not to emit. Add a loom test for it.
  3. Replace NLJ's two counters, the incomplete flag and the unused per-chunk counter with ProbeCompletion.

Tests to run:

  • the existing NLJ and hash join unit tests, joins.slt, nested_loop_join_spill.slt and join_fuzz
  • RUSTFLAGS="--cfg datafusion_loom" cargo test -p datafusion-physical-plan --lib loom_tests. CI does not run this, so please paste the result in the PR.

Describe alternatives you've considered

No response

Additional context

Possible follow-ups:

#25413 also changes nested_loop_join.rs.

Activity

  1. changed the title [-]Joins: one shared protocol for "the last probe partition emits the build-side rows"[/-] [+]Use `ProbeCompletion` in `NestedLoopJoinExec`[/+] on Oct 9, 2026
  2. jeffoodchain commented on Oct 9, 2026

    @jeffoodchain

    @jayzhan211 I will be working on this.

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or requesthelp wantedExtra attention is needed

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions