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]