jayzhan211 opened a new issue, #26148:
URL: https://github.com/apache/datafusion/issues/26148

   ### Is your feature request related to a problem or challenge?
   
   Several joins build one side once and probe it from many partitions: 
`HashJoinExec` (`CollectLeft`), `NestedLoopJoinExec` and 
`PiecewiseMergeJoinExec`. They all follow the same rule:
   
   > Every probe partition records which build rows it matched. When the 
**last** partition finishes, it emits the build rows the join type needs at the 
end, and no other partition does. For example, it emits unmatched rows for 
LEFT/FULL joins and matched rows for LEFT SEMI joins.
   
   Each operator implements this rule itself. On main (55b2f093bb) there are 
four copies, using three different memory orderings:
   
   | Where | Counter | Ordering | How the last partition gets the matches |
   |---|---|---|---|
   | `HashJoinExec` | `ProbeCompletion` (`hash_join/probe_completion.rs`), 
model-checked with loom | AcqRel | `finish_cloned()` on the shared bitmap |
   | `NestedLoopJoinExec`, in memory | `JoinLeftData::probe_threads_counter` 
(`nested_loop_join.rs:1147`, `:1181`) | Relaxed | locks the shared bitmap |
   | `NestedLoopJoinExec`, spilled left side | 
`LeftSpillData::{probe_threads_counter, incomplete}` and `depart_unfinished` 
(`:1463-1580`) | AcqRel / Release | takes the bitmap |
   | `PiecewiseMergeJoinExec` | `remaining_partitions` 
(`piecewise_merge_join/exec.rs:1063`), decremented in `classic_join.rs:227` and 
`existence_join.rs:284` | SeqCst | `min_marked` watermark |
   
   There is more duplication around these copies:
   
   - The NLJ spill path builds a `JoinLeftData` for each chunk with a counter 
of `1` that, per its comment, "is not used" (`nested_loop_join.rs:2563-2571`).
   - Five functions answer "does this join type emit build rows at the end?":
     - `need_produce_result_in_final`;
     - a PWMJ function with the same name and a different answer;
     - `need_to_produce_result_in_final` in SHJ;
     - `emits_unmatched_left_rows`;
     - `need_produce_right_in_final`.
   
   So anyone changing one copy first has to work out two things for that copy 
alone:
   - **Why its memory ordering is correct.** For example, `Relaxed` in NLJ is 
safe only because the bitmap sits behind a `Mutex`.
   - **What happens when a partition is dropped early.**
   
   Each copy was fixed separately this year:
   
   - #22791: NLJ decremented twice and emitted spurious unmatched rows. #22865 
then added a `ProbeEnd` state just to enforce "decrement once".
   - #24675: the per-chunk `JoinLeftData` in the NLJ fallback had a counter of 
1. Every partition thought it was last, and the join returned wrong rows.
   - #25004: a dropped NLJ partition never reported, so the other partitions 
hung.
   - #25076: HJ read the NULL flag before decrementing and returned wrong `NOT 
IN` results. This fix added `ProbeCompletion` and its loom model, but only HJ 
uses them.
   - #25542: rebuilt the NLJ fallback state, with its own counter and 
`incomplete` flag.
   
   ### Describe the solution you'd like
   
   One concept and one type:
   
   > A shared build side owns one **completion** object. Each probe partition 
calls it exactly once: it reports that it finished, or that it went away early. 
The call that brings the count to zero is told it is last. It emits the build 
rows only if no partition went away early. Either way it releases the shared 
state. Operators never touch the counter or the memory orderings directly.
   
   For `HashJoinExec`, `ProbeCompletion` is already this type, and it has a 
loom model. This issue proposes making it the only implementation, in small 
steps.
   
   **First PR (this issue): use `ProbeCompletion` in `NestedLoopJoinExec`**
   
   1. Move `hash_join/probe_completion.rs` to `joins/probe_completion.rs` as 
`pub(crate)`. HJ only changes its import.
   2. Add a call for a partition that goes away unfinished, e.g. 
`abandon(count)`. The last reporter then knows that nobody should emit, but can 
still release shared state. Add this case to the loom model. The one copy of 
the orderings has to stay under loom.
   3. In NLJ, use one `ProbeCompletion` per left side, both in memory and 
spilled. Remove `JoinLeftData::probe_threads_counter`, including the per-chunk 
dummy, and `LeftSpillData::{probe_threads_counter, incomplete}`. 
`depart_unfinished` becomes the abandon call.
   
   The NOT IN facts that HJ records (`saw_row`, `saw_null_key`) can stay on the 
shared type, since NLJ never records them, or they can move to an HJ-side 
wrapper. The implementer can choose.
   
   **Expected size:** about 200-250 changed lines, counting the file move as a 
rename. Production code should shrink slightly.
   
   **Expected behaviour change:** none.
   - HJ is untouched apart from the import.
   - The in-memory decrement in NLJ goes from `Relaxed` to `AcqRel`. That costs 
one fence per partition, not per row.
   
   **Acceptance criteria**
   
   - `cargo test -p datafusion-physical-plan --lib joins::nested_loop_join` and 
`joins::hash_join` pass. The NLJ tests that assert on `probe_threads_counter` 
are adapted to the new type.
   - `RUSTFLAGS="--cfg datafusion_loom" cargo test -p datafusion-physical-plan 
--lib loom_tests` passes, including a new case for abandon. CI does not run 
loom, so include the result in the PR description.
   - The NLJ spill tests in `core/tests/memory_limit`, `join_fuzz`, `joins.slt` 
(which contains the #22791 reproducer) and `nested_loop_join_spill.slt` pass.
   - `nested_loop_join.rs` no longer declares any atomics of its own for probe 
completion.
   
   **Later steps** (separate issues, each a small PR)
   
   - `PiecewiseMergeJoinExec` uses the same type, after #25840 lands.
   - Bundle the shared match bitmap with the completion object, so the last 
partition receives the finished bitmap by value. This removes HJ's 
`finish_cloned()` and NLJ's lock per emitted range.
   - Replace the five predicates with one helper.
   - With a single type in place, #12454 (a hook to share the match state for 
distributed `CollectLeft`) has one obvious place to attach.
   
   ### Describe alternatives you've considered
   
   - **Leave it as it is.** New work such as HJ spilling (#24768), NLJ 
semi/anti pruning (#25413) and #12454 each touch one copy and can reintroduce 
one of the bugs above.
   - **Also unify the per-operator stream state machines.** Too 
operator-specific, and not PR-sized.
   
   ### Additional context
   
   - **Coordination:** #25413 (NLJ semi/anti pruning) is open on the same file, 
and whichever PR lands second rebases. HJ logic is not changed, so this does 
not conflict with the hash join PRs in flight.
   - **Related:** #12454, #21650 (where shared join state should live), #24768.
   


-- 
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]

Reply via email to