fix: cancel the NestedLoopJoin coordinated fallback when a partition is dropped unfinished - #25004
Open
viirya wants to merge 10 commits into
Open
fix: cancel the NestedLoopJoin coordinated fallback when a partition is dropped unfinished#25004viirya wants to merge 10 commits into
viirya wants to merge 10 commits into
Conversation
…ished The partitions of a coordinated memory-limited fallback are not independent: chunk advancement needs every one of them to report. A partition that goes away before finishing therefore never reports, so no emitter is elected, nothing releases the coordinator slot, and the survivors fall through to the `notified()` wait in `next_chunk` and hang. Because the coordinator is owned by the exec rather than the streams, it also keeps holding the chunk it published, so its reservation stays charged for as long as the plan is alive. Reproduced directly against the coordinator: one partition takes chunk 0 and disappears, and a second partition asking for chunk 1 never returns. Three changes, which only make sense together: 1. The coordinator lock becomes a synchronous `parking_lot::Mutex`. Every critical section here was already synchronous -- the one slow operation, `load_one_chunk`, runs after the guard is dropped -- so `next_chunk` is restructured to decide under the lock, release it, then act: a `Decision` enum carries `Serve` / `Finished` / `Load` / `Wait` / `Cancelled` out of the locked block. This is what lets cleanup complete synchronously instead of depending on a future that a dropped stream takes with it. 2. A `cancelled` flag plus `cancel()`, driven from `Drop for NestedLoopJoinStream`. A stream that reached `Done` finished its work and cancels nothing; any other state cancels the fallback, dropping `current`, `carryover`, the shared left stream and the coordinator's reservation, and waking the waiters so they return an error rather than blocking. The guard exists from stream construction, so it also covers a cancellation that happens before the stream ever takes a chunk. Chunks other partitions still hold stay accounted until they release them. 3. A loader that was reading outside the lock must not publish into a coordinator that `cancel` has already cleaned up; doing so would reinstate the stream and reservation it just dropped. The publish path now discards its result and reports the cancellation instead. This accepts a deliberate semantic: cancelling one unfinished partition cancels the coordinated execution. It can no longer produce a complete result, so failing the remaining partitions is more honest than leaving them hanging or letting them report success from partial input. Making the release synchronous also removes the machinery that existed only to drive it across polls: `chunk_release_in_flight`, its three poll sites, `handle_releasing_final_chunk`, the `NLJState::ReleasingFinalChunk` state, and the `cx` argument `handle_emit_left_unmatched` no longer needs. With no future to interrupt, the "release dropped while pending" shape disappears on its own. Tests cover the cancellation timings that matter: a dropped partition not hanging the survivors, cancel releasing what the coordinator held, cancellation while a load is in flight discarding its result, and a normally-completed stream not cancelling anything. Each fails if `cancel()` is neutered -- the hang test by timing out, the memory test with 505 bytes still reserved. Co-authored-by: Claude Code
Review of the previous commit found three places where cancellation was not actually observed, plus a hole in my own test coverage. Drop only handled `SpillState::Active`, so a stream cancelled while still `Pending` -- set up for the fallback but not yet holding a chunk -- reported nothing and cancelled nothing. That is exactly the "cancelled before taking a chunk" case the guard was supposed to cover, so the claim that placing it at stream construction was sufficient was wrong. Both states now cancel; only `Disabled` does not. Cancellation was only checked inside `next_chunk`, which a stream holding the final chunk never calls again. Such a survivor ran to completion and reported success built from an execution that had lost a partition. `poll_next` now checks on every iteration. `cancel()` could not interrupt a load already in flight: the loader awaits `load_one_chunk` directly, and waking `notify` does not reach a read parked on its input. The load now races the cancellation signal, and on cancellation drops its local stream and reservation, clears the leader claim and reports the cancellation rather than waiting for a read that may never finish. The publish check alone prevented reinstatement but not the stall. The test gap is the more important lesson. My four tests all called `coordinator.cancel()` directly, so neutering `cancel()` failed them while neutering the entire `Drop` body left all 68 tests passing -- they exercised the function but never the wiring that calls it. The four review tests added here drive real plans and streams, and two of them fail when `Drop` is neutered. Co-authored-by: Claude Code
…iters Follow-up review found two remaining hang paths, both traceable to using the chunk-progress `Notify` to carry cancellation. `Notify` broadcasts to waiters that already exist and stores no state, so a cancellation landing before a waiter registered was simply lost, and a task parked on something other than chunk progress had no waker registered at all. Cancellation now travels on its own `cancel_notify`, and every waiter registers with `Notified::enable()` before reading the `cancelled` flag. A cancellation arriving between those two steps is delivered rather than lost, which closes the gap at initial registration. Because `cancel_notify` is signalled only by `cancel()`, a wake from it always means a real cancellation: the re-arm loop that re-registered after unrelated wakes is gone, and with it the race that loop carried. For streams, checking a flag was not enough. A stream returning Pending from its right input is not waiting on the coordinator, so dropping a peer never woke it and the poll-loop check was not reached until some unrelated event happened to poll the task. Streams now hold a registered `cancellation_watcher` across polls and poll it each iteration, so the waker really is with the coordinator. Writing a replacement for the now-inapplicable re-arm test caught a third bug of my own: the loader's watcher was still constructed from the progress `Notify`, so any chunk-progress traffic resolved it and failed the load as if cancelled. The new test drives unrelated `notify_waiters()` past a parked loader and requires it to stay pending, then cancels for real. The obsolete re-arm test is dropped -- it drove a loop that no longer exists -- and the replacement needs no instrumentation. The hook for cancelling before watcher registration stays, since that gap still needs guarding. Verified: 8 review tests pass; joins 1162; memory_limit 38 including apache#24746's regressions; clippy clean. Neutering `Drop` fails 4 of the review tests, and weakening the stream check back to a plain flag read fails 1, so both mechanisms are covered rather than merely present. Co-authored-by: Claude Code
Round 3 found no blocking correctness issue: the reviewer walked the four notification interleavings and confirmed that enabling the `Notified` before reading `cancelled` covers each one, and that a single-signaller broadcast makes the removed re-arm loop unnecessary. What it did find was documentation and test names claiming more than the code delivers. Adds the reviewer's four tests: a watcher constructed before cancellation but first polled after it, plus one constructed after; several live watchers with waker replacement and progress traffic that must not wake them; a real `NestedLoopJoinStream` parked on build input receiving a cancellation wakeup when its peer is dropped; and both partitions running to normal completion and being dropped, which must leave the coordinator usable, return all nine rows, and free the memory while the plan stays alive. That last one closes the gap my own "completed stream" test never covered. Documentation corrections, all cases of promising more than happens: - `cancellation_watcher` said the future was already queued. Registration happens on its first poll, so callers have to poll it, not just hold it. - `cancel_notify` was described as notified "once, permanently". It is a broadcast carrying no state; `cancelled` is what persists. - `cancel`'s documentation had been split by a later insertion and was sitting above `is_cancelled`. Moved back, and it now records that `cancel` is the only signaller of `cancel_notify` -- the invariant the missing re-arm loop rests on. Two of my tests were renamed because their names overstated what they did. `test_nlj_cancel_during_load_discards_the_result` cancels *before* starting the load, so it exercises the entry check, not a publish after an in-flight read; it is now `test_nlj_cancelled_coordinator_refuses_to_serve_chunks` and points at the test that does pause a real read. `test_nlj_completed_stream_drop_does_not_cancel` never built or dropped a stream; it is now `test_nlj_uncancelled_coordinator_serves_and_stays_live` and points at the reviewer's test that drops real streams. Verified: joins 1166; memory_limit 38 including apache#24746's regressions; clippy clean. Neutering `Drop` fails 5 review tests and weakening the stream check to a plain flag read fails 2. Co-authored-by: Claude Code
Round 3 listed these as optional, sharing the already-reviewed watcher and Drop mechanisms. They are cheap and they close the two cases where I was reasoning rather than testing. A stream can park in `EmitGlobalRightUnmatched` rather than `BufferingLeft`: that state reopens the spilled right side and polls it. Since the watcher is polled at the top of every `poll_next` iteration regardless of state, it should be woken there too, and now that is asserted rather than assumed -- a FULL join parked on a pending replay records a wake when its peer is dropped. An unfinished stream that ends in an error also reaches `Drop` without passing through `Done`, so it cancels its peers. That follows from the `Done` check, but the error path had no test of its own; one now injects a failing right input and confirms the peers are cancelled. Both fail if `Drop` is neutered. Dropping `Pending` back out of `Drop`'s match -- the round-1 defect -- now fails 5 review tests rather than the 1 it did then. Verified: 14 review tests; joins 1168; memory_limit 38 including apache#24746's regressions; clippy clean. While adding these I removed a `#[tokio::test]` that belonged to `review_v3_cancel_wakes_pending_build_input` and restored it, then audited every test in the module for a missing attribute. Only the shared `review_paused_loader` helper lacks one, correctly. This is the second time an attribute has gone missing to careless editing here, and a silently unregistered test is worse than a failing one. Co-authored-by: Claude Code
Review round 4 found that neither fixture added in the previous commit exercised its stated scenario, and both were rewritten rather than patched. The error fixture built a stream with a failing right input and dropped it without polling, so it only repeated "an unstarted stream cancels its peers", which other tests already cover. Polling alone would not have helped: its `left_data` was a permanently pending future, so it could never leave `BufferingLeft` to reach the right input at all. It now resolves the build side through `LeftLoad::Spilled`, drives the stream under a bounded timeout until the injected error surfaces, and asserts that error and a non-`Done` state *before* dropping -- so a fixture that stops reaching the right input fails instead of quietly degrading. The replay fixture assigned `EmitGlobalRightUnmatched` on top of `SpillState::Pending`, a combination execution never produces. It reached the pending read only because `right_data` was already `Some`, bypassing the `Active`-only reopen branch, so it was really testing a `Pending` watcher while claiming to test the replay configuration. It now polls until the spilled build side puts it in `Active`, asserts that, and only then parks in the replay stage. The replay reader is still injected rather than reopened -- that shortcut is now stated in the doc comment -- and the test additionally requires the wake to surface the cancellation rather than merely counting a wakeup. Also corrects the rustdoc above `test_nlj_cancelled_coordinator_refuses_to_serve_chunks`, which still described cancelling mid-load and validating publication after the rename fixed the name and the body comment, and softens the `cancel_notify` wording: a delivered broadcast is itself sufficient, so it is not merely a nudge to go and read the flag. On test registration: `cargo test -- --list` reports all 14 review tests, which is the right check. The regex sweep I ran was the wrong tool -- the module has legitimate async helpers, so requiring an attribute on every `async fn` would be wrong -- and it is not committed. Verified: 14 review tests; joins 1168; memory_limit 38 including apache#24746's regressions; clippy clean. Neutering `Drop` fails both fixtures; removing the watcher poll fails the replay one only, which is the expected split since the error fixture does not depend on the watcher. Co-authored-by: Claude Code
Review round 5 reported no blocking findings and cleared the change for upstream, but pointed at three comments still describing the pre-refactor code. The coordinator's inner state was documented as guarded by an async mutex; it is a `parking_lot::Mutex`, and the reason matters enough to record: cancellation and chunk release have to finish inside one `poll_next` rather than depending on a future a dropped stream would take away, and no critical section awaits. The coordinator reservation was described as holding the current chunk's memory, registered via `initiate_fallback` -- a function that no longer exists. It only holds bytes while a load runs; `load_one_chunk` then moves them into the chunk's `JoinLeftData` with `take()`, so accounting follows the data. The comment above `buffered_left_data = None` claimed the `Arc` reaches zero once the last partition lets go. The slot holds a strong reference too, so the reservation is freed only when every holder *and* the slot release. That is precisely why the slot cannot hold the chunk weakly, which two earlier attempts here established the hard way, so the comment now says it. Also corrects an overreach in my own reporting: I had described the sqllogictest suite as blocked by the pre-existing `needless_pass_by_value` failure in the sqllogictest crate. That was wrong -- clippy escalates warnings with `-D warnings`, `cargo test` does not. Running it works: `cargo test -p datafusion-sqllogictest --test sqllogictests -- nested_loop_join joins information_schema` passes 8/8 here. The lint failure is real but separate, and reproduces on clean `upstream/main` (`3266eaa91`) with `cargo clippy -p datafusion-sqllogictest --all-targets --all-features -- -D warnings`. Verified: joins 1168; sqllogictest 8/8 on the join and information_schema files; clippy clean for datafusion-physical-plan. Co-authored-by: Claude Code
Round 6 cleared the change for upstream and corrected three pieces of wording. The important one: I had justified the slot's strong `Arc` by saying the reservation frees only after every holder and the slot release. True, but that is a consequence, not the reason. The reason is a handoff: a faster partition can drop its reference before a slower one has taken the chunk at all, and the slow one is served from the slot (`Decision::Serve`), so the slot has to keep the chunk alive across that gap. Stated that way, the comment explains why holding it weakly cannot work -- which two earlier attempts here established by failing. "Complete inside a single `poll_next`" was wrong about cancellation, which runs from `Drop` and has neither a poll nor an await; it now says "without another poll or await". "Only holds bytes while a load is in progress" was too absolute: the error path returns the reservation to the coordinator with its bytes still accounted, since `take()` happens only on success. The comment now describes the successful path without claiming the stronger invariant. Verified: joins 1168; clippy clean for datafusion-physical-plan. Comment-only, so memory_limit and sqllogictest were not rerun. Co-authored-by: Claude Code
Review-process naming and duplication that should not reach upstream. The fourteen tests carried `review_v2`/`v3`/`v4` prefixes recording which review round produced them, which means nothing to anyone reading the file later. They are renamed after the behaviour they check, under one `nlj_` prefix: what is woken, what is released, what refuses to serve. The two shared helpers are renamed the same way (`run_paused_loader_cancellation`, `cancellation_test_plan`). The coordinator carried a test-only `review_cancel_before_watch` field with no documentation at all -- a field named after a review round, sitting in production state. It is now `cancel_at_leader_claim` and says what it is for: firing a cancellation in the window between claiming a load and registering the watcher, which is the interleaving where a lost notification would strand every other partition, and which a test cannot otherwise reach. Two `eprintln!`s left over from debugging are gone; the assertions beside them already carried the meaning. The `WakeCount` waker was defined identically in four tests and is now one helper with `new`/`count`/`reset` -- deliberately just that, not a test framework. While doing the extraction a regex of mine rewrote the helper's own body into `self.reset()` calling itself, which the test run caught as a stack overflow rather than a wrong answer. Fixed, and a reminder that mechanical renames need the result read back. Verified: 44 NLJ tests, all registered per `cargo test -- --list`; joins 1168; clippy clean. Neutering `Drop` still fails 7 of them, and disabling the renamed test seam fails the test that depends on it, so the renames did not quietly detach any coverage. Co-authored-by: Claude Code
Round 7 found no blocking issue but caught five precision problems, one of them mine from the previous commit. Inserting the shared `WakeCount` put it directly beneath an existing rustdoc, so it captured documentation belonging to `multi_partition_memory_limited_join_collect_concurrent` -- the block explaining why that helper must collect partitions concurrently rather than sequentially. `WakeCount` now carries only its own docs and the explanation is back on the helper it describes. Same class of mistake as the recursive-body regex last commit: an insertion that looked local but moved something adjacent. `nlj_normal_completion_leaves_coordinator_usable` claimed more than it checks -- it establishes that normal completion does not cancel peers of the same execution, not that a plan or coordinator can be re-executed (it cannot; the coordinator's state is one-shot). Renamed to `nlj_normal_completion_does_not_cancel_peers`. `nlj_drop_while_pending_releases_chunk` named a chunk release that has not happened at that point: no chunk is loaded before the drop. It asserts the memory is not retained afterwards, so it is now `nlj_drop_pending_partition_does_not_retain_chunk_memory`. Removing `Pending` from `Drop`'s match still fails it, so the name change tracks the same coverage. The test seam's documentation said it reproduced an interleaving "a test cannot otherwise hit" and that a lost cancellation would "strand every other partition". Both overstated: it reproduces the interleaving deterministically, and what it strands is the loader itself, since other observers may already be returning errors. Verified: joins 1168; clippy clean. For the record, the nested_loop_join module holds 82 tests; the `nlj_` prefix selects 44 of them, so earlier reports of "44 NLJ tests" were describing the filter, not the module. Co-authored-by: Claude Code
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #25004 +/- ##
==========================================
+ Coverage 81.64% 81.71% +0.06%
==========================================
Files 1124 1126 +2
Lines 413173 415774 +2601
Branches 413173 415774 +2601
==========================================
+ Hits 337320 339729 +2409
- Misses 55995 56064 +69
- Partials 19858 19981 +123 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Which issue does this PR close?
Rationale for this change
In the coordinated memory-limited NestedLoopJoin fallback, dropping an unfinished output partition can strand the shared chunk and leave other partitions waiting indefinitely. This change cancels the coordinated execution synchronously, releases coordinator-owned resources, and wakes unfinished partitions and in-flight loaders so they report an execution error. Normally completed partitions do not cancel their peers, and chunks still held by streams remain memory-accounted until released.
The timing that hangs: chunk advancement needs every probing partition to report, since the last one to report becomes the emitter and releases the coordinator slot. A partition that vanishes never reports, so no emitter is elected and nothing releases the slot. A survivor then asks for the next chunk and finds Case 1 unmatched (the slot still holds the previous chunk) while Cases 2 and 3 both require
current.is_none(), so it falls into thenotified()wait innext_chunkwith nothing left to wake it. Because the coordinator is owned by the exec rather than the streams, a plan that outlives them also keeps the slot'sArc<JoinLeftData>, and with it the chunk's reservation.Two shapes reach this: a release future dropped while pending on the coordinator mutex, and — more broadly — a mid-probe cancellation where no emitter is ever elected, so
release_chunkis never called at all. The second is why hardening the release future alone cannot fix it.What changes are included in this PR?
The coordinator lock is now synchronous (
parking_lot::Mutex). Every critical section already was: the one slow operation,load_one_chunk, runs after the guard is dropped.next_chunkis restructured to decide under the lock and act after releasing it, carrying aDecisionout of the locked block. This is what lets cleanup finish without another poll or await — cancellation runs fromDrop, which has neither.Cancellation travels on its own broadcast with registered waiters.
cancel()sets a permanentcancelledflag and signals a dedicatedcancel_notify; it is the only signaller, so a wake from it always means a real cancellation. Waiters enable theirNotifiedbefore reading the flag, so a cancellation landing between those steps is delivered rather than lost. Streams hold a registeredcancellation_watcheracross polls and poll it every iteration — a stream parked on its right input is not waiting on chunk progress, so a flag read alone would never reach it.A loader reading outside the lock cannot publish into a cancelled coordinator. The publish path would otherwise reinstate the stream and reservation that
cancelhad just dropped. The load also races the cancellation signal, so a read parked on input does not have to finish before the cancellation is observed.Dead machinery removed. With the release synchronous,
chunk_release_in_flight, its poll sites,handle_releasing_final_chunkandNLJState::ReleasingFinalChunkhad no remaining purpose, and the "release dropped while pending" shape disappears with them.What is the testing strategy for this PR?
Fourteen tests in
nested_loop_join.rscovering the cancellation timings: a dropped partition not hanging its peers; cancellation releasing coordinator-held memory; cancellation before watcher registration; chunk-progress traffic not resolving a cancellation watcher; cancellation during an in-flight load discarding its publish, and not waiting for a read that may never finish; wakeups for streams parked on build input, right input, and the global-right replay; an errored stream cancelling peers on drop; and normal completion of both partitions leaving the coordinator alone, returning all rows and freeing memory while the plan stays alive.They are mutation-checked rather than merely passing: neutering
Dropfails 7 of them, removingPendingfromDrop's match fails the retention test, and weakening the stream-side check to a plain flag read fails the parked-stream tests.Also run:
joins(1168),memory_limit(38, including #24746's regressions), thenested_loop_join/joins/information_schemasqllogictest files (8/8), andcargo clippy -p datafusion-physical-plan --all-targets --all-features -- -D warnings.Unrelated and pre-existing:
cargo clippy -p datafusion-sqllogictest --all-targets --all-features -- -D warningsfailsneedless_pass_by_valueatdatafusion/sqllogictest/src/engines/conversion.rs:99on cleanmain(3266eaa91) as well.Are there any user-facing changes?
Dropping an unfinished partition cancels the coordinated NLJ execution, causing other unfinished partitions to report an execution error. Normally completed partitions do not cancel their peers. Chunks remain memory-accounted while any stream still references them.
Previously those partitions hung, so this replaces an indefinite wait with a reported failure. Reviewers may want to weigh in on that semantic specifically: the alternative would be letting survivors finish, which cannot produce a complete result once a partition is gone.
The existing opt-out behavior is unchanged: it disables coordinated fallback for multi-partition joins requiring final left-side emission. No API changes.