You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
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:
#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:
In memory:JoinLeftData::probe_threads_counter uses Relaxed. That is correct only because the bitmap happens to sit behind a Mutex, so every reader has to work that out again.
Spill path:LeftSpillData::{probe_threads_counter, incomplete} and depart_unfinished add a separate flag for point 3.
Low risk: it is about 200 lines with no intended behaviour change. Going from Relaxed to AcqRel costs one fence per partition, not per row.
Describe the solution you'd like
Use ProbeCompletion in NestedLoopJoinExec too:
Move hash_join/probe_completion.rs to joins/probe_completion.rs.
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.
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.
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
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:
JoinLeftDatawith a counter of1. 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.Relaxedcounter, soNOT INreturned 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:report_completed(), so they cannot be read before the decrement.AcqRel.NestedLoopJoinExecstill provides all four guarantees by hand, in two places with different orderings:JoinLeftData::probe_threads_counterusesRelaxed. That is correct only because the bitmap happens to sit behind aMutex, so every reader has to work that out again.LeftSpillData::{probe_threads_counter, incomplete}anddepart_unfinishedadd a separate flag for point 3.JoinLeftDatawhose counter is always1and, per its comment, "is not used". That is the shape that caused fix: refuse memory-limited NestedLoopJoin fallback for left-emission joins with a multi-partition probe side #24675.Moving NLJ onto
ProbeCompletionbrings:CollectLeft(Proposal: Hook to better supportCollectLeftjoins in distributed execution #12454).RelaxedtoAcqRelcosts one fence per partition, not per row.Describe the solution you'd like
Use
ProbeCompletioninNestedLoopJoinExectoo:hash_join/probe_completion.rstojoins/probe_completion.rs.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.incompleteflag and the unused per-chunk counter withProbeCompletion.Tests to run:
joins.slt,nested_loop_join_spill.sltandjoin_fuzzRUSTFLAGS="--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:
PiecewiseMergeJoinExec, which has a third copy of this counter (SeqCst). Wait until perf(pwmj):settle never-matching streamed rows against the buffered extreme… #25840 lands.#25413 also changes
nested_loop_join.rs.