This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/main/pr-25004-41f5ea0e02dbce204f443c05017a1dcfdffe8ff1
in repository https://gitbox.apache.org/repos/asf/datafusion.git

commit f8ecabbf4bafbefdebbe789e5b3ac18b6b02bb0d
Author: Liang-Chi Hsieh <[email protected]>
AuthorDate: Wed Sep 9 15:28:48 2026 +0000

    fix: cancel the NestedLoopJoin coordinated fallback when a partition is 
dropped unfinished (#25004)
    
    ## Which issue does this PR close?
    
    - Closes #25003.
    
    ## 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 the `notified()` wait in `next_chunk` with 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's
    `Arc<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_chunk` is 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_chunk` is
    restructured to decide under the lock and act after releasing it,
    carrying a `Decision` out of the locked block. This is what lets cleanup
    finish without another poll or await — cancellation runs from `Drop`,
    which has neither.
    
    **Cancellation travels on its own broadcast with registered waiters.**
    `cancel()` sets a permanent `cancelled` flag and signals a dedicated
    `cancel_notify`; it is the only signaller, so a wake from it always
    means a real cancellation. Waiters enable their `Notified` before
    reading the flag, so a cancellation landing between those steps is
    delivered rather than lost. Streams hold a registered
    `cancellation_watcher` across 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 `cancel` had 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_chunk` and `NLJState::ReleasingFinalChunk` had
    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.rs` covering 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 `Drop`
    fails 7 of them, removing `Pending` from `Drop`'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), the `nested_loop_join` / `joins` / `information_schema`
    sqllogictest files (8/8), and `cargo 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 warnings` fails
    `needless_pass_by_value` at
    `datafusion/sqllogictest/src/engines/conversion.rs:99` on clean `main`
    (`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.
---
 .../physical-plan/src/joins/nested_loop_join.rs    | 1875 ++++++++++++++++----
 1 file changed, 1578 insertions(+), 297 deletions(-)

diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs 
b/datafusion/physical-plan/src/joins/nested_loop_join.rs
index 9f0abd92a0..2fe48df474 100644
--- a/datafusion/physical-plan/src/joins/nested_loop_join.rs
+++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs
@@ -65,8 +65,8 @@ use datafusion_common::cast::as_boolean_array;
 use datafusion_common::tree_node::TreeNodeRecursion;
 use datafusion_common::{
     JoinSide, NullEquality, Result, ScalarValue, Statistics, arrow_err,
-    assert_eq_or_internal_err, internal_datafusion_err, internal_err, 
project_schema,
-    unwrap_or_internal_err,
+    assert_eq_or_internal_err, exec_err, internal_datafusion_err, internal_err,
+    project_schema, unwrap_or_internal_err,
 };
 use datafusion_execution::memory_pool::{MemoryConsumer, MemoryReservation};
 use datafusion_execution::{SpillFile, TaskContext};
@@ -1341,17 +1341,6 @@ enum NLJState {
     /// has to guard against decrementing twice.
     ProbeEnd,
     EmitLeftUnmatched,
-    /// Drives the final chunk's `release_chunk` future to completion.
-    ///
-    /// Non-final chunks get released on the way back through
-    /// `BufferingLeft`, but after the last chunk the stream goes straight to
-    /// `Done` or `EmitGlobalRightUnmatched`, and neither polls
-    /// `chunk_release_in_flight`. Since the coordinator hangs off the exec
-    /// rather than off this stream, leaving that future unpolled keeps the
-    /// final chunk's batch, bitmap and reservation accounted for as long as
-    /// the plan is alive. This state exists to poll it exactly once, then
-    /// continue to whichever state would have followed.
-    ReleasingFinalChunk,
     /// Emit unmatched right rows using the global bitmap accumulated across
     /// all left chunks. Only used in memory-limited mode for join types that
     /// require tracking right-side matches in the final output (RIGHT, FULL,
@@ -1400,12 +1389,20 @@ struct CurrentChunk {
     is_last: bool,
 }
 
-/// Inner state of [`FallbackCoordinator`], guarded by an async mutex.
+/// Inner state of [`FallbackCoordinator`], guarded by a synchronous mutex.
+///
+/// Synchronous because cancellation and chunk release have to complete without
+/// another poll or await -- cancellation runs from `Drop`, which has neither 
--
+/// rather than depending on a future a dropped stream would take with it. No
+/// critical section awaits: the one slow operation, reading a chunk, runs 
after
+/// the guard is released.
 struct FallbackCoordinatorInner {
-    /// Reservation owned by the coordinator. Holds the memory for the
-    /// currently-loaded chunk. Reset (`resize(0)`) between chunks.
-    /// Lazily registered by the first leader, after the runtime context
-    /// becomes available via `initiate_fallback`.
+    /// Reservation the leader borrows to bound one chunk load.
+    ///
+    /// On a successful load `load_one_chunk` moves the accounted bytes into 
the
+    /// chunk's `JoinLeftData` with `take()`, so the accounting follows the 
data
+    /// rather than staying with this slot. Lazily registered by the first
+    /// leader, once a runtime context is available.
     reservation: Option<MemoryReservation>,
     /// The shared left spill stream from which chunks are read. Owned by
     /// the coordinator so only one partition reads it at a time.
@@ -1428,6 +1425,36 @@ struct FallbackCoordinatorInner {
     /// True while a partition has claimed leader role for the next
     /// chunk and is loading it; prevents two partitions from racing.
     loader_in_flight: bool,
+    /// A partition was dropped while still `Pending`, before the shared load
+    /// decided whether the left side spills.
+    ///
+    /// `Pending` is entered by every eligible execution, including those whose
+    /// left side ends up fitting in memory -- and those never build a shared
+    /// chunk counter, so their right partitions stay independent and a dropped
+    /// peer has nothing to coordinate. Cancelling on such a drop would fail a
+    /// query that had no fallback at all, so it is only recorded here.
+    ///
+    /// Paired with `coordination_started`: whichever of the two happens second
+    /// performs the cancellation, so the drop is honoured whether it precedes 
or
+    /// follows the execution becoming coordinated. A bool suffices -- one lost
+    /// partition is enough to cancel, and nothing reads a count.
+    pending_drop: bool,
+    /// Set once any partition has entered the coordinated path.
+    ///
+    /// Remembered rather than checked in the moment, because a `Pending` drop 
can
+    /// arrive after a peer is already coordinating; that drop must cancel, and
+    /// without this flag there would be nobody left to notice.
+    coordination_started: bool,
+    /// Set when a stream is dropped before finishing, which cancels the whole
+    /// coordinated fallback.
+    ///
+    /// The partitions of a coordinated fallback are not independent: chunk
+    /// advancement requires every one of them to report, so a partition that
+    /// disappears mid-probe would otherwise leave the survivors waiting on a
+    /// release nobody will ever make. Once set, chunk state is dropped, 
waiters
+    /// are woken with an error, and a loader that is still reading must 
discard
+    /// its result instead of publishing it.
+    cancelled: bool,
 }
 
 /// Plan-level shared coordinator for the memory-limited fallback path.
@@ -1444,10 +1471,31 @@ pub(crate) struct FallbackCoordinator {
     /// Whether `JoinLeftData` should carry a left visited bitmap (for
     /// join types that emit unmatched left rows in the final output).
     with_visited_bitmap: bool,
-    inner: tokio::sync::Mutex<FallbackCoordinatorInner>,
+    inner: Mutex<FallbackCoordinatorInner>,
     /// Notified when a new chunk becomes available, when the left stream
     /// is exhausted, or when a chunk is released.
     notify: tokio::sync::Notify,
+    /// Broadcast signalled when the fallback is cancelled.
+    ///
+    /// This carries no state of its own -- `cancelled` is what persists. A
+    /// delivered broadcast is itself sufficient to establish cancellation;
+    /// observers read the flag before awaiting, so neither alone is relied on.
+    /// Kept separate from `notify`
+    /// because cancellation has to reach tasks that are not waiting on chunk
+    /// progress at all: a stream parked on its right input, or a loader parked
+    /// on a spill read. Waiters enable their `Notified` before reading
+    /// `cancelled`, so a cancellation landing between those two steps is
+    /// delivered rather than lost.
+    cancel_notify: tokio::sync::Notify,
+    /// Test seam that reproduces one cancellation interleaving 
deterministically.
+    ///
+    /// Set to `1` to make the leader cancel after claiming the load but before
+    /// it registers its cancellation watcher. A cancellation lost in that 
window
+    /// strands the loader itself -- other observers may already be returning
+    /// errors -- which is what the paired test checks. Consumed when it fires,
+    /// so one store arms it once.
+    #[cfg(test)]
+    cancel_at_leader_claim: AtomicUsize,
 }
 
 impl FallbackCoordinator {
@@ -1455,7 +1503,7 @@ impl FallbackCoordinator {
         Self {
             right_partition_count,
             with_visited_bitmap,
-            inner: tokio::sync::Mutex::new(FallbackCoordinatorInner {
+            inner: Mutex::new(FallbackCoordinatorInner {
                 reservation: None,
                 left_stream: None,
                 left_schema: None,
@@ -1464,28 +1512,145 @@ impl FallbackCoordinator {
                 next_chunk_index: 0,
                 current: None,
                 loader_in_flight: false,
+                pending_drop: false,
+                coordination_started: false,
+                cancelled: false,
             }),
             notify: tokio::sync::Notify::new(),
+            cancel_notify: tokio::sync::Notify::new(),
+            #[cfg(test)]
+            cancel_at_leader_claim: AtomicUsize::new(0),
         }
     }
 
     /// After the last partition finishes processing chunk
     /// `released_chunk_index`, drop the slot so the next leader can
     /// load chunk `released_chunk_index + 1`.
-    async fn release_chunk(self: &Arc<Self>, released_chunk_index: usize) {
-        let mut inner = self.inner.lock().await;
-        if let Some(cur) = &inner.current
-            && cur.chunk_index == released_chunk_index
+    fn release_chunk(self: &Arc<Self>, released_chunk_index: usize) {
         {
-            inner.current = None;
-            inner.next_chunk_index = released_chunk_index + 1;
+            let mut inner = self.inner.lock();
+            if let Some(cur) = &inner.current
+                && cur.chunk_index == released_chunk_index
+            {
+                inner.current = None;
+                inner.next_chunk_index = released_chunk_index + 1;
+            }
         }
         // Always notify: waiters may be blocked because they couldn't
         // become leader while a previous chunk was current.
-        drop(inner);
         self.notify.notify_waiters();
     }
 
+    /// True once a partition was dropped unfinished, cancelling the fallback.
+    ///
+    /// Production code observes cancellation through `cancellation_watcher`, 
so
+    /// that a waker is registered; this plain read is for assertions only.
+    #[cfg(test)]
+    fn is_cancelled(&self) -> bool {
+        self.inner.lock().cancelled
+    }
+
+    /// A future that resolves when the fallback is cancelled.
+    ///
+    /// Registration happens on the future's **first poll**, not at 
construction:
+    /// that poll enables the `Notified` and then reads `cancelled`, so 
whichever
+    /// happens first is observed. Callers must therefore poll it, not merely 
hold
+    /// it.
+    ///
+    /// Callers that park on something other than chunk progress -- a stream
+    /// waiting on its right input, for instance -- need one of these polled
+    /// alongside their own work, or a peer's cancellation never reaches their
+    /// waker.
+    fn cancellation_watcher(self: &Arc<Self>) -> BoxFuture<'static, ()> {
+        let coordinator = Arc::clone(self);
+        async move {
+            let notified = coordinator.cancel_notify.notified();
+            let mut notified = std::pin::pin!(notified);
+            notified.as_mut().enable();
+            if coordinator.inner.lock().cancelled {
+                return;
+            }
+            notified.await;
+        }
+        .boxed()
+    }
+
+    /// Records a partition dropped before the shared load decided whether the
+    /// left side spills.
+    ///
+    /// Whether this cancels depends on what the surviving partitions are 
doing,
+    /// which is why the drop is recorded rather than acted on unconditionally:
+    ///
+    /// * No peer has coordinated yet -- only `pending_drop` is set. If the 
load
+    ///   resolves to `InMemory` nobody ever coordinates and this stays inert, 
so
+    ///   an execution that never needed the coordinator is not failed by it.
+    /// * A peer is already coordinating -- cancel now. That peer is waiting 
on a
+    ///   probe report this partition will never make.
+    ///
+    /// The second case is the mirror of [`Self::begin_coordination`]: the two
+    /// share `pending_drop` and `coordination_started` under one lock, so
+    /// whichever runs second performs the cancellation and neither order is 
lost.
+    fn record_pending_drop(self: &Arc<Self>) {
+        let cancel_now = {
+            let mut inner = self.inner.lock();
+            if inner.cancelled {
+                return;
+            }
+            inner.pending_drop = true;
+            // A peer may already be coordinating, in which case this drop is 
not
+            // hypothetical: that peer is waiting on a report this partition 
will
+            // never make.
+            inner.coordination_started
+        };
+        if cancel_now {
+            self.cancel();
+        }
+    }
+
+    /// Records that this execution is now coordinating, and honours any drop 
that
+    /// happened while it was not.
+    ///
+    /// Called when a stream enters memory-limited mode. Setting the flag 
matters
+    /// as much as the check: a `Pending` partition dropped *after* this point 
has
+    /// to cancel too, and `record_pending_drop` reads this flag to decide 
that.
+    fn begin_coordination(self: &Arc<Self>) {
+        let cancel_now = {
+            let mut inner = self.inner.lock();
+            inner.coordination_started = true;
+            !inner.cancelled && inner.pending_drop
+        };
+        if cancel_now {
+            self.cancel();
+        }
+    }
+
+    /// Cancels the coordinated fallback and drops everything the coordinator
+    /// holds, synchronously.
+    ///
+    /// Called from a stream's drop guard when it goes away without finishing.
+    /// Chunks other partitions still hold stay accounted until they release
+    /// them; what this drops is the coordinator's own state, which nothing 
will
+    /// come back for.
+    ///
+    /// This is the only place that signals `cancel_notify`. Watchers rely on
+    /// that: a wake from it always means a real cancellation, which is why 
they
+    /// need no re-arm loop. Keep it that way if you add call sites.
+    fn cancel(self: &Arc<Self>) {
+        {
+            let mut inner = self.inner.lock();
+            if inner.cancelled {
+                return;
+            }
+            inner.cancelled = true;
+            inner.current = None;
+            inner.carryover = None;
+            inner.left_stream = None;
+            inner.reservation = None;
+        }
+        self.notify.notify_waiters();
+        self.cancel_notify.notify_waiters();
+    }
+
     /// Fetch `expected_chunk_index`, becoming leader to load it from the
     /// left spill stream if no other partition has done so. Returns
     /// `Ok(None)` when the left stream is exhausted and no chunk with
@@ -1500,135 +1665,218 @@ impl FallbackCoordinator {
         // `spill_data` is already resolved: every partition receives the
         // same `Arc<LeftSpillData>` from the shared `OnceAsync<LeftLoad>`,
         // so the left child is executed and spilled exactly once.
-
         loop {
-            let mut inner = self.inner.lock().await;
-
-            // Case 1: requested chunk is already loaded.
-            if let Some(cur) = &inner.current
-                && cur.chunk_index == expected_chunk_index
-            {
-                return Ok(Some((Arc::clone(&cur.data), cur.is_last)));
-            }
-
-            // Case 2: left stream exhausted and no current chunk to
-            // deliver — caller is past the last chunk.
-            if inner.left_exhausted
-                && inner.current.is_none()
-                && inner.carryover.is_none()
-            {
-                return Ok(None);
-            }
+            // Decide what to do with the lock held, then act on that decision
+            // after releasing it. The guard must not survive into the `.await`
+            // below -- it is not `Send`, and holding it across the load would
+            // serialize every partition behind the leader's disk reads.
+            let decision = {
+                let mut inner = self.inner.lock();
+
+                if inner.cancelled {
+                    Decision::Cancelled
+                } else if let Some(cur) = &inner.current
+                    && cur.chunk_index == expected_chunk_index
+                {
+                    // Case 1: requested chunk is already loaded.
+                    Decision::Serve(Arc::clone(&cur.data), cur.is_last)
+                } else if inner.left_exhausted
+                    && inner.current.is_none()
+                    && inner.carryover.is_none()
+                {
+                    // Case 2: left side finished and nothing left to deliver.
+                    Decision::Finished
+                } else if inner.current.is_none() && !inner.loader_in_flight {
+                    // Case 3: claim the leader role and take the shared
+                    // resources out so the load can run without the lock.
+                    inner.loader_in_flight = true;
+                    let stream = inner.left_stream.take();
+                    let reservation = inner.reservation.take();
+                    let carryover = inner.carryover.take();
+                    let chunk_index_to_load = inner.next_chunk_index;
+                    debug_assert_eq!(chunk_index_to_load, 
expected_chunk_index);
+                    Decision::Load {
+                        stream,
+                        reservation,
+                        carryover,
+                        chunk_index: chunk_index_to_load,
+                    }
+                } else {
+                    // Case 4: someone else is loading, or this chunk index has
+                    // already been passed -- wait to be notified.
+                    Decision::Wait(self.notify.notified())
+                }
+            };
 
-            // Case 3: no chunk loaded and no leader yet — claim leader.
-            if inner.current.is_none() && !inner.loader_in_flight {
-                inner.loader_in_flight = true;
-                // Lazily initialize the shared left stream, schema, and
-                // chunk reservation. Only the leader does this (under the
-                // `loader_in_flight` guard), so waiters never open a
-                // throwaway spill stream that the leader would overwrite.
-                let mut left_stream = match inner.left_stream.take() {
-                    Some(stream) => stream,
-                    None => {
-                        // Construct the spill stream. If this ever fails,
-                        // clear the leader flag and wake waiters before
-                        // returning, so they don't block forever on a
-                        // release that the failed leader will never make.
-                        match spill_data.spill_manager.read_spill_as_stream(
-                            Arc::clone(&spill_data.spill_file),
-                            None,
-                        ) {
-                            Ok(stream) => {
-                                inner.left_schema = 
Some(Arc::clone(&spill_data.schema));
-                                stream
-                            }
-                            Err(e) => {
-                                inner.loader_in_flight = false;
-                                drop(inner);
-                                self.notify.notify_waiters();
-                                return Err(e);
+            match decision {
+                Decision::Cancelled => {
+                    return exec_err!(
+                        "NestedLoopJoin coordinated fallback was cancelled 
because a \
+                         partition was dropped before finishing"
+                    );
+                }
+                Decision::Serve(data, is_last) => return Ok(Some((data, 
is_last))),
+                Decision::Finished => return Ok(None),
+                Decision::Wait(notified) => {
+                    notified.await;
+                }
+                Decision::Load {
+                    stream,
+                    reservation,
+                    carryover,
+                    chunk_index,
+                } => {
+                    // Build whatever the slot did not already have. A failure
+                    // here must clear the leader flag and wake waiters, or 
they
+                    // block on a release the failed leader never makes.
+                    let (mut left_stream, left_schema) = match stream {
+                        Some(stream) => {
+                            let schema = {
+                                let inner = self.inner.lock();
+                                inner.left_schema.clone()
+                            };
+                            let schema = match schema {
+                                Some(schema) => schema,
+                                None => Arc::clone(&spill_data.schema),
+                            };
+                            (stream, schema)
+                        }
+                        None => {
+                            match 
spill_data.spill_manager.read_spill_as_stream(
+                                Arc::clone(&spill_data.spill_file),
+                                None,
+                            ) {
+                                Ok(stream) => {
+                                    let mut inner = self.inner.lock();
+                                    inner.left_schema =
+                                        Some(Arc::clone(&spill_data.schema));
+                                    drop(inner);
+                                    (stream, Arc::clone(&spill_data.schema))
+                                }
+                                Err(e) => {
+                                    {
+                                        let mut inner = self.inner.lock();
+                                        inner.loader_in_flight = false;
+                                    }
+                                    self.notify.notify_waiters();
+                                    return Err(e);
+                                }
                             }
                         }
+                    };
+                    let mut reservation = match reservation {
+                        Some(reservation) => reservation,
+                        None => {
+                            
MemoryConsumer::new("NestedLoopJoinFallbackChunk".to_string())
+                                .with_can_spill(true)
+                                .register(task_context.memory_pool())
+                        }
+                    };
+
+                    // Race the read against cancellation. A loader parked on
+                    // its input is not waiting on `notify`, so without this a
+                    // `cancel` would not be observed until the read finished 
on
+                    // its own -- which may be never if the input is gone.
+                    // Test seam: fire a cancellation in the window between
+                    // claiming the load and registering the watcher below. See
+                    // `cancel_at_leader_claim`.
+                    #[cfg(test)]
+                    if self.cancel_at_leader_claim.swap(0, Ordering::SeqCst) 
== 1 {
+                        self.cancel();
                     }
-                };
-                let mut reservation = match inner.reservation.take() {
-                    Some(reservation) => reservation,
-                    None => {
-                        
MemoryConsumer::new("NestedLoopJoinFallbackChunk".to_string())
-                            .with_can_spill(true)
-                            .register(task_context.memory_pool())
-                    }
-                };
-                let left_schema = Arc::clone(
-                    inner
-                        .left_schema
-                        .as_ref()
-                        .expect("left_schema installed above"),
-                );
-                let carryover = inner.carryover.take();
-                let chunk_index_to_load = inner.next_chunk_index;
-                debug_assert_eq!(chunk_index_to_load, expected_chunk_index);
-                drop(inner);
-
-                let load_result = Arc::clone(&self)
-                    .load_one_chunk(
-                        chunk_index_to_load,
+                    let cancelled = self.cancel_notify.notified();
+                    let load = Arc::clone(&self).load_one_chunk(
+                        chunk_index,
                         &mut left_stream,
                         &mut reservation,
                         carryover,
                         Arc::clone(&left_schema),
                         build_time.clone(),
-                    )
-                    .await;
-
-                // Re-acquire lock and publish the result.
-                let mut inner = self.inner.lock().await;
-                inner.left_stream = Some(left_stream);
-                inner.reservation = Some(reservation);
-                inner.loader_in_flight = false;
-
-                match load_result {
-                    Ok(LoadOutcome::Chunk {
-                        data,
-                        is_last,
-                        carryover,
-                    }) => {
-                        inner.carryover = carryover;
-                        if is_last {
-                            inner.left_exhausted = true;
+                    );
+                    let load_result = {
+                        let mut load = std::pin::pin!(load);
+                        let mut cancelled = std::pin::pin!(cancelled);
+                        // Queue the waiter, then read the flag. If 
cancellation
+                        // already happened we bail without awaiting; if it 
lands
+                        // just after, the enabled future receives the 
broadcast
+                        // rather than losing it.
+                        cancelled.as_mut().enable();
+                        if self.inner.lock().cancelled {
+                            None
+                        } else {
+                            tokio::select! {
+                                biased;
+                                result = &mut load => Some(result),
+                                () = &mut cancelled => None,
+                            }
+                        }
+                    };
+                    let Some(load_result) = load_result else {
+                        // Cancelled mid-read. Drop the local stream and
+                        // reservation instead of putting them back, and
+                        // clear the leader claim so nothing waits on us.
+                        drop(left_stream);
+                        drop(reservation);
+                        {
+                            let mut inner = self.inner.lock();
+                            inner.loader_in_flight = false;
                         }
-                        let arc_data = Arc::new(data);
-                        inner.current = Some(CurrentChunk {
-                            chunk_index: chunk_index_to_load,
-                            data: Arc::clone(&arc_data),
-                            is_last,
-                        });
-                        drop(inner);
-                        self.notify.notify_waiters();
-                        return Ok(Some((arc_data, is_last)));
-                    }
-                    Ok(LoadOutcome::Empty) => {
-                        // No data at all. Mark exhausted; let other
-                        // partitions observe and exit.
-                        inner.left_exhausted = true;
-                        drop(inner);
-                        self.notify.notify_waiters();
-                        return Ok(None);
-                    }
-                    Err(e) => {
-                        drop(inner);
                         self.notify.notify_waiters();
-                        return Err(e);
+                        return exec_err!(
+                            "NestedLoopJoin coordinated fallback was cancelled 
\
+                             while a chunk was being loaded"
+                        );
+                    };
+
+                    // Publish, unless the fallback was cancelled while this
+                    // load was running -- putting the stream and reservation
+                    // back would undo the cleanup `cancel` just did.
+                    let published = {
+                        let mut inner = self.inner.lock();
+                        inner.loader_in_flight = false;
+                        if inner.cancelled {
+                            None
+                        } else {
+                            inner.left_stream = Some(left_stream);
+                            inner.reservation = Some(reservation);
+                            match load_result {
+                                Ok(LoadOutcome::Chunk {
+                                    data,
+                                    is_last,
+                                    carryover,
+                                }) => {
+                                    inner.carryover = carryover;
+                                    if is_last {
+                                        inner.left_exhausted = true;
+                                    }
+                                    let arc_data = Arc::new(data);
+                                    inner.current = Some(CurrentChunk {
+                                        chunk_index,
+                                        data: Arc::clone(&arc_data),
+                                        is_last,
+                                    });
+                                    Some(Ok(Some((arc_data, is_last))))
+                                }
+                                Ok(LoadOutcome::Empty) => {
+                                    inner.left_exhausted = true;
+                                    Some(Ok(None))
+                                }
+                                Err(e) => Some(Err(e)),
+                            }
+                        }
+                    };
+                    self.notify.notify_waiters();
+                    match published {
+                        Some(result) => return result,
+                        None => {
+                            return exec_err!(
+                                "NestedLoopJoin coordinated fallback was 
cancelled \
+                                 while a chunk was being loaded"
+                            );
+                        }
                     }
                 }
             }
-
-            // Case 4: another partition is loading the next chunk, or
-            // the current chunk is for a previous index we've already
-            // moved past — wait to be notified.
-            let notified = self.notify.notified();
-            drop(inner);
-            notified.await;
         }
     }
 
@@ -1739,6 +1987,26 @@ impl FallbackCoordinator {
     }
 }
 
+/// What [`FallbackCoordinator::next_chunk`] decided to do while holding the
+/// lock, so the work itself can happen after the guard is released.
+enum Decision<'a> {
+    /// Chunk is loaded and can be handed straight back.
+    Serve(Arc<JoinLeftData>, bool),
+    /// Left side is finished; the caller is past the last chunk.
+    Finished,
+    /// This caller is the leader and owns the shared resources for one load.
+    Load {
+        stream: Option<SendableRecordBatchStream>,
+        reservation: Option<MemoryReservation>,
+        carryover: Option<RecordBatch>,
+        chunk_index: usize,
+    },
+    /// Someone else is loading, or this index has been passed: wait.
+    Wait(tokio::sync::futures::Notified<'a>),
+    /// A partition was dropped before finishing, so the fallback is cancelled.
+    Cancelled,
+}
+
 enum LoadOutcome {
     Chunk {
         data: JoinLeftData,
@@ -1824,9 +2092,6 @@ pub(crate) struct SpillStateActive {
     /// until it resolves to either the next chunk or `None` (left side
     /// exhausted with no chunk to deliver).
     chunk_fetch_in_flight: Option<ChunkFetchFuture>,
-    /// In-flight chunk release future. Created by `EmitLeftUnmatched`
-    /// when the last partition for a chunk has finished its work.
-    chunk_release_in_flight: Option<BoxFuture<'static, ()>>,
 }
 
 impl SpillStateActive {
@@ -1908,6 +2173,22 @@ pub(crate) struct NestedLoopJoinStream {
     // ========================================================================
     /// State Tracking
     state: NLJState,
+    /// Set when the stream terminated by reporting a cancellation.
+    ///
+    /// `Done` normally means "finished cleanly", and `handle_done` pads an 
empty
+    /// result with one empty batch so an all-filtered join still carries its
+    /// schema. A cancelled stream has already returned an error and must not 
then
+    /// emit anything, empty batch included, so it takes the terminal path
+    /// directly.
+    cancelled_terminally: bool,
+    /// Registered watcher for the coordinator's cancellation broadcast.
+    ///
+    /// Polled on every `poll_next` iteration so a stream parked on its own 
input
+    /// still has a waker registered with the coordinator; without it a peer's
+    /// cancellation would not be observed until some unrelated event happened 
to
+    /// wake this task.
+    cancel_watch: Option<BoxFuture<'static, ()>>,
+
     /// Output buffer holds the join result to output. It will emit eagerly 
when
     /// the threshold is reached.
     output_buffer: Box<BatchCoalescer>,
@@ -2014,6 +2295,42 @@ impl Stream for NestedLoopJoinStream {
         cx: &mut std::task::Context<'_>,
     ) -> Poll<Option<Self::Item>> {
         loop {
+            // A peer partition may have been dropped unfinished at any point,
+            // including while this stream was working on the final chunk. The
+            // coordinated execution cannot produce a complete result after 
that,
+            // so fail rather than emit output from partial input.
+            if !matches!(self.state, NLJState::Done) {
+                // Poll a registered watcher rather than only reading the flag:
+                // this stream may be about to park on its own input, and the
+                // poll leaves a waker with the coordinator so a peer's
+                // cancellation actually reaches it.
+                if self.cancel_watch.is_none() {
+                    self.cancel_watch = match &self.spill_state {
+                        SpillState::Active(active) => {
+                            Some(active.coordinator.cancellation_watcher())
+                        }
+                        SpillState::Pending {
+                            fallback_coordinator,
+                            ..
+                        } => Some(fallback_coordinator.cancellation_watcher()),
+                        SpillState::Disabled => None,
+                    };
+                }
+                let cancelled = match self.cancel_watch.as_mut() {
+                    Some(watch) => watch.poll_unpin(cx).is_ready(),
+                    None => false,
+                };
+                if cancelled {
+                    self.state = NLJState::Done;
+                    self.cancelled_terminally = true;
+                    self.cancel_watch = None;
+                    return Poll::Ready(Some(exec_err!(
+                        "NestedLoopJoin coordinated fallback was cancelled 
because \
+                         a partition was dropped before finishing"
+                    )));
+                }
+            }
+
             match self.state {
                 // # NLJState transitions
                 // --> FetchingRight
@@ -2155,17 +2472,7 @@ impl Stream for NestedLoopJoinStream {
                     let join_metric = 
self.metrics.join_metrics.join_time.clone();
                     let _join_timer = join_metric.timer();
 
-                    match self.handle_emit_left_unmatched(cx) {
-                        ControlFlow::Continue(()) => {}
-                        ControlFlow::Break(poll) => {
-                            return 
self.metrics.join_metrics.baseline.record_poll(poll);
-                        }
-                    }
-                }
-
-                // Release the final chunk before finishing the stream.
-                NLJState::ReleasingFinalChunk => {
-                    match self.handle_releasing_final_chunk(cx) {
+                    match self.handle_emit_left_unmatched() {
                         ControlFlow::Continue(()) => {}
                         ControlFlow::Break(poll) => {
                             return 
self.metrics.join_metrics.baseline.record_poll(poll);
@@ -2213,6 +2520,45 @@ impl RecordBatchStream for NestedLoopJoinStream {
     }
 }
 
+/// Cancels the coordinated fallback if a stream goes away before finishing.
+///
+/// The partitions of a coordinated fallback are not independent: chunk
+/// advancement needs every one of them to report, so a partition that
+/// disappears mid-probe would leave the survivors waiting on a release nobody
+/// will ever make, and the coordinator -- owned by the plan, not the stream --
+/// would keep holding the chunk it had published. Cancelling is the honest
+/// outcome: this execution can no longer produce a complete result, so the
+/// remaining partitions are failed rather than left hanging or allowed to
+/// report success from partial input.
+///
+/// A stream that reached `Done` finished its work and does not cancel 
anything.
+impl Drop for NestedLoopJoinStream {
+    fn drop(&mut self) {
+        if matches!(self.state, NLJState::Done) {
+            return;
+        }
+        // `Active` means this stream was probing shared chunks, so its peers 
are
+        // already relying on coordination and have to be told now.
+        //
+        // `Pending` is different: it is entered before the shared load decides
+        // whether the left side spills at all, so it does not yet imply a
+        // coordinated execution. Record the drop instead and let the shared
+        // `coordination_started` flag decide: if a peer already coordinates, 
the
+        // recording call itself cancels; if none has yet, the next one to 
enter
+        // the coordinated path does. An execution whose left side fits in 
memory
+        // never enters it and so never cancels -- which is the point, since it
+        // has no coordination to lose.
+        match &self.spill_state {
+            SpillState::Active(active) => 
Arc::clone(&active.coordinator).cancel(),
+            SpillState::Pending {
+                fallback_coordinator,
+                ..
+            } => Arc::clone(fallback_coordinator).record_pending_drop(),
+            SpillState::Disabled => {}
+        }
+    }
+}
+
 impl NestedLoopJoinStream {
     #[expect(clippy::too_many_arguments)]
     pub(crate) fn new(
@@ -2240,6 +2586,8 @@ impl NestedLoopJoinStream {
             current_right_batch: None,
             current_right_batch_matched: None,
             state: NLJState::BufferingLeft,
+            cancelled_terminally: false,
+            cancel_watch: None,
             left_probe_idx: 0,
             left_emit_idx: 0,
             left_exhausted: false,
@@ -2314,7 +2662,6 @@ impl NestedLoopJoinStream {
             global_right_bitmaps_reservation,
             right_batch_index: 0,
             chunk_fetch_in_flight: None,
-            chunk_release_in_flight: None,
         }));
 
         // State stays BufferingLeft — next poll will enter
@@ -2353,7 +2700,15 @@ impl NestedLoopJoinStream {
                              entering memory-limited mode"
                         );
                         match 
self.enter_memory_limited_mode(Arc::clone(left_spill)) {
-                            Ok(()) => ControlFlow::Continue(()),
+                            Ok(()) => {
+                                // Now that this execution is known to 
coordinate,
+                                // a partition dropped while still `Pending`
+                                // becomes a cancellation.
+                                if let SpillState::Active(active) = 
&self.spill_state {
+                                    
Arc::clone(&active.coordinator).begin_coordination();
+                                }
+                                ControlFlow::Continue(())
+                            }
                             Err(e) => 
ControlFlow::Break(Poll::Ready(Some(Err(e)))),
                         }
                     }
@@ -2379,18 +2734,6 @@ impl NestedLoopJoinStream {
             );
         };
 
-        // Drain any pending chunk-release future before fetching the next
-        // chunk. The coordinator slot must be released so the next leader
-        // can load the following chunk.
-        if let Some(fut) = active.chunk_release_in_flight.as_mut() {
-            match fut.poll_unpin(cx) {
-                Poll::Ready(()) => {
-                    active.chunk_release_in_flight = None;
-                }
-                Poll::Pending => return ControlFlow::Break(Poll::Pending),
-            }
-        }
-
         // Lazily start a chunk-fetch future for `active.next_chunk_index`.
         if active.chunk_fetch_in_flight.is_none() {
             let coordinator = Arc::clone(&active.coordinator);
@@ -2658,27 +3001,12 @@ impl NestedLoopJoinStream {
     /// next chunk (if the left stream is not yet exhausted).
     fn handle_emit_left_unmatched(
         &mut self,
-        cx: &mut std::task::Context<'_>,
     ) -> ControlFlow<Poll<Option<Result<RecordBatch>>>> {
         // Return any completed batches first
         if let Some(poll) = self.maybe_flush_ready_batch() {
             return ControlFlow::Break(poll);
         }
 
-        // First, drive any pending chunk-release future to completion so
-        // we don't transition to the next state while another partition
-        // is waiting for the slot to be freed.
-        if let SpillState::Active(active) = &mut self.spill_state
-            && let Some(fut) = active.chunk_release_in_flight.as_mut()
-        {
-            match fut.poll_unpin(cx) {
-                Poll::Ready(()) => {
-                    active.chunk_release_in_flight = None;
-                }
-                Poll::Pending => return ControlFlow::Break(Poll::Pending),
-            }
-        }
-
         // Process current unmatched state
         match self.process_left_unmatched() {
             // State unchanged (EmitLeftUnmatched)
@@ -2696,9 +3024,17 @@ impl NestedLoopJoinStream {
                     }
 
                     // Drop our reference to the current chunk's
-                    // `JoinLeftData` before releasing the slot. Once the
-                    // last partition does this, the `Arc` reaches zero
-                    // refcount and the per-chunk reservation is freed.
+                    // `JoinLeftData` before releasing the slot. The 
coordinator
+                    // slot holds a strong `Arc` too, so the reservation inside
+                    // the chunk is freed only once every probing partition 
*and*
+                    // the slot have let go.
+                    //
+                    // The slot's reference is load-bearing, not redundant: a
+                    // faster partition can reach here and let go before a 
slower
+                    // one has taken the chunk at all, and the slow one is 
served
+                    // from the slot (`Decision::Serve`). The slot therefore 
has
+                    // to keep the chunk alive across that handoff, which is 
why
+                    // it cannot hold it weakly.
                     self.buffered_left_data = None;
 
                     if self.is_memory_limited() {
@@ -2709,14 +3045,11 @@ impl NestedLoopJoinStream {
                             // releases the coordinator slot so the next
                             // leader can load the following chunk.
                             if is_emitter {
+                                // Synchronous now, so the slot is freed before
+                                // this poll returns. It can no longer be lost 
by
+                                // the stream being dropped mid-release.
                                 let coordinator = 
Arc::clone(&active.coordinator);
-                                let released_index = active.next_chunk_index;
-                                active.chunk_release_in_flight = Some(
-                                    async move {
-                                        
coordinator.release_chunk(released_index).await
-                                    }
-                                    .boxed(),
-                                );
+                                
coordinator.release_chunk(active.next_chunk_index);
                             }
                             active.next_chunk_index += 1;
                         }
@@ -2725,18 +3058,7 @@ impl NestedLoopJoinStream {
                         // does not need to be reset here.
                     }
 
-                    if self.is_memory_limited()
-                        && self.left_exhausted
-                        && matches!(
-                            &self.spill_state,
-                            SpillState::Active(active)
-                                if active.chunk_release_in_flight.is_some()
-                        )
-                    {
-                        // Final chunk: hand off to `ReleasingFinalChunk`, 
which
-                        // drives the release future before moving on.
-                        self.state = NLJState::ReleasingFinalChunk;
-                    } else if !self.left_exhausted && self.is_memory_limited() 
{
+                    if !self.left_exhausted && self.is_memory_limited() {
                         // More left data to process — go back to
                         // BufferingLeft for the next chunk.
                         self.left_probe_idx = 0;
@@ -2762,39 +3084,6 @@ impl NestedLoopJoinStream {
         }
     }
 
-    /// Handle ReleasingFinalChunk state.
-    ///
-    /// Polls the final chunk's `release_chunk` future to completion, then
-    /// moves on to whichever state would have followed `EmitLeftUnmatched`.
-    /// Releasing the slot drops the coordinator's `Arc<JoinLeftData>` and lets
-    /// the coordinator reservation be reclaimed, which otherwise would not
-    /// happen until the whole plan is dropped.
-    fn handle_releasing_final_chunk(
-        &mut self,
-        cx: &mut std::task::Context<'_>,
-    ) -> ControlFlow<Poll<Option<Result<RecordBatch>>>> {
-        if let SpillState::Active(active) = &mut self.spill_state
-            && let Some(fut) = active.chunk_release_in_flight.as_mut()
-        {
-            match fut.poll_unpin(cx) {
-                Poll::Ready(()) => {
-                    active.chunk_release_in_flight = None;
-                }
-                Poll::Pending => return ControlFlow::Break(Poll::Pending),
-            }
-        }
-
-        self.state = if self.should_track_unmatched_right {
-            // Drop the exhausted right stream so that
-            // EmitGlobalRightUnmatched opens a fresh replay pass.
-            self.right_data = None;
-            NLJState::EmitGlobalRightUnmatched
-        } else {
-            NLJState::Done
-        };
-        ControlFlow::Continue(())
-    }
-
     /// Handle EmitGlobalRightUnmatched state.
     ///
     /// Replays all right batches from the spill file and emits unmatched
@@ -2902,6 +3191,12 @@ impl NestedLoopJoinStream {
 
     /// Handle Done state - final state processing
     fn handle_done(&mut self) -> Poll<Option<Result<RecordBatch>>> {
+        // A cancelled stream already reported its error. Nothing may follow 
it --
+        // not buffered output, and not the empty-result padding below.
+        if self.cancelled_terminally {
+            return Poll::Ready(None);
+        }
+
         // Return any remaining completed batches before final termination
         if let Some(poll) = self.maybe_flush_ready_batch() {
             return poll;
@@ -4258,7 +4553,7 @@ pub(crate) mod tests {
                     
MemoryConsumer::new("NestedLoopJoinFallbackChunk[test]".to_string())
                         .register(task_ctx.memory_pool()),
                 ));
-                let mut inner = coordinator.inner.lock().await;
+                let mut inner = coordinator.inner.lock();
                 inner.left_exhausted = true;
                 inner.current = Some(CurrentChunk {
                     chunk_index: 0,
@@ -4266,7 +4561,7 @@ pub(crate) mod tests {
                     is_last: true,
                 });
             } else {
-                let mut inner = coordinator.inner.lock().await;
+                let mut inner = coordinator.inner.lock();
                 inner.left_schema = Some(Arc::clone(&left_schema));
                 inner.left_stream = Some(left_stream);
             }
@@ -4296,7 +4591,6 @@ pub(crate) mod tests {
                 global_right_bitmaps_reservation,
                 right_batch_index: 0,
                 chunk_fetch_in_flight: None,
-                chunk_release_in_flight: None,
             };
             (
                 OnceFut::new(async { internal_err!("unused left data was 
polled") }),
@@ -5660,13 +5954,32 @@ pub(crate) mod tests {
         )
     }
 
-    /// Run a NLJ across 4 right partitions, collecting every output
-    /// partition CONCURRENTLY. This is required for the multi-chunk
-    /// coordinator path: a chunk is not released until all partitions
-    /// finish probing it, and a partition cannot advance to the next chunk
-    /// until the current one is released. Collecting partitions
-    /// sequentially would therefore deadlock; concurrent collection mirrors
-    /// how partitions actually run under the runtime.
+    /// Waker that counts how often it is woken.
+    ///
+    /// Cancellation tests assert that a parked stream is *woken*, not merely 
that
+    /// a flag flipped, so they need to see the waker fire.
+    struct WakeCount(AtomicUsize);
+
+    impl futures::task::ArcWake for WakeCount {
+        fn wake_by_ref(this: &Arc<Self>) {
+            this.0.fetch_add(1, Ordering::SeqCst);
+        }
+    }
+
+    impl WakeCount {
+        fn new() -> Arc<Self> {
+            Arc::new(Self(AtomicUsize::new(0)))
+        }
+
+        fn count(&self) -> usize {
+            self.0.load(Ordering::SeqCst)
+        }
+
+        fn reset(&self) {
+            self.0.store(0, Ordering::SeqCst);
+        }
+    }
+
     /// Spills a plan's single partition to a `LeftSpillData`, so a test can 
hand
     /// the coordinator the same input the fallback path would.
     async fn spill_left_for_test(
@@ -5698,66 +6011,1027 @@ pub(crate) mod tests {
         }))
     }
 
-    /// The chunk's memory must be owned by the chunk's `JoinLeftData`, so it
-    /// stays accounted while *any* holder still references it.
-    ///
-    /// The emitter elected in `ProbeEnd` is the last stream to finish probing,
-    /// which is not necessarily the last to drop its `Arc<JoinLeftData>`: a
-    /// non-emitter can flush a completed output batch from
-    /// `maybe_flush_ready_batch` and return while still holding one, since 
that
-    /// return happens before `buffered_left_data = None`. Freeing the bytes 
when
-    /// the coordinator slot is released would therefore under-account memory
-    /// that is still live.
-    ///
-    /// This drives the coordinator directly so the check does not depend on
-    /// scheduling: after releasing the slot, the pool must still account for 
the
-    /// chunk while a reference is held, and drop to zero only once it is gone.
     #[tokio::test]
-    async fn test_nlj_chunk_memory_is_owned_by_the_chunk_data() -> Result<()> {
+    async fn nlj_cancel_wakes_stream_parked_on_build_input() -> Result<()> {
         let runtime = RuntimeEnvBuilder::new().build_arc()?;
-        let pool = Arc::clone(&runtime.memory_pool);
-        let task_ctx = Arc::new(TaskContext::default().with_runtime(runtime));
-
-        // One chunk, tracked bitmap, two nominal probe partitions.
+        let ctx = Arc::new(TaskContext::default().with_runtime(runtime));
         let coordinator = Arc::new(FallbackCoordinator::new(2, true));
-        let left = build_left_table();
-        let spill = spill_left_for_test(Arc::clone(&left), 
Arc::clone(&task_ctx)).await?;
-
-        let (chunk, is_last) = Arc::clone(&coordinator)
-            .next_chunk(0, Arc::clone(&spill), Arc::clone(&task_ctx), 
Time::new())
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+        let _chunk = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&ctx), Time::new())
             .await?
-            .expect("the left side has rows, so a chunk must be produced");
-        assert!(is_last, "the fixture fits in a single chunk");
-
-        let accounted_while_held = pool.reserved();
+            .expect("chunk");
+        let make_stream = || {
+            let right_schema = build_right_table().schema();
+            let (schema, columns) =
+                build_join_schema(&spill.schema, &right_schema, 
&JoinType::Left);
+            NestedLoopJoinStream::new(
+                Arc::new(schema),
+                None,
+                JoinType::Left,
+                Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                    right_schema,
+                    futures::stream::pending::<Result<RecordBatch>>(),
+                )),
+                OnceFut::new(futures::future::pending::<Result<LeftLoad>>()),
+                columns,
+                NestedLoopJoinMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
+                1,
+                SpillState::Pending {
+                    task_context: Arc::clone(&ctx),
+                    fallback_coordinator: Arc::clone(&coordinator),
+                },
+            )
+        };
+        let peer = make_stream();
+        let mut survivor = make_stream();
+        let wakes = WakeCount::new();
+        let waker = futures::task::waker(Arc::clone(&wakes));
+        let mut cx = std::task::Context::from_waker(&waker);
+        assert!(survivor.poll_next_unpin(&mut cx).is_pending());
+        assert!(matches!(survivor.state, NLJState::BufferingLeft));
+        wakes.reset();
+        drop(peer);
+        // Both streams are parked on a build load that never resolves, so 
neither
+        // ever enters the coordinated path and the drop is only recorded -- 
which
+        // is the correct behaviour, since an execution whose left side might 
still
+        // fit in memory has no coordination to lose.
         assert!(
-            accounted_while_held > 0,
-            "loading a chunk must account for its batch and bitmap"
+            !coordinator.is_cancelled(),
+            "a pending drop must not cancel before anything coordinates"
         );
-
-        // Release the coordinator slot while still holding the chunk. This is
-        // the emitter's release: the slot is freed, but the data is alive.
-        coordinator.release_chunk(0).await;
-        assert_eq!(
-            pool.reserved(),
-            accounted_while_held,
-            "releasing the slot must not release memory that is still 
referenced"
+        // Once some partition does begin coordinating, the recorded drop 
applies
+        // and has to reach this parked stream.
+        Arc::clone(&coordinator).begin_coordination();
+        assert!(coordinator.is_cancelled());
+        assert!(
+            wakes.count() > 0,
+            "cancel must wake the survivor waiting on build input"
         );
+        Ok(())
+    }
 
-        // Dropping the last reference is what returns the bytes.
-        drop(chunk);
-        assert_eq!(
-            pool.reserved(),
-            0,
-            "dropping the last chunk reference must release its memory"
+    #[tokio::test]
+    async fn nlj_cancel_observed_when_watcher_polled_after_cancellation() {
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let before = coordinator.cancellation_watcher();
+        coordinator.cancel();
+        coordinator.cancel();
+        let after = coordinator.cancellation_watcher();
+        assert!(before.now_or_never().is_some());
+        assert!(after.now_or_never().is_some());
+    }
+
+    #[tokio::test]
+    async fn nlj_cancel_reaches_every_watcher_and_survives_waker_replacement() 
{
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let mut first = coordinator.cancellation_watcher();
+        let mut second = coordinator.cancellation_watcher();
+        let old_count = WakeCount::new();
+        let new_count = WakeCount::new();
+        let other_count = WakeCount::new();
+        let old_waker = futures::task::waker(Arc::clone(&old_count));
+        let new_waker = futures::task::waker(Arc::clone(&new_count));
+        let other_waker = futures::task::waker(Arc::clone(&other_count));
+        assert!(
+            first
+                .poll_unpin(&mut std::task::Context::from_waker(&old_waker))
+                .is_pending()
         );
-        Ok(())
+        assert!(
+            first
+                .poll_unpin(&mut std::task::Context::from_waker(&new_waker))
+                .is_pending()
+        );
+        assert!(
+            second
+                .poll_unpin(&mut std::task::Context::from_waker(&other_waker))
+                .is_pending()
+        );
+        coordinator.notify.notify_waiters();
+        assert_eq!(new_count.count(), 0);
+        assert_eq!(other_count.count(), 0);
+        coordinator.cancel();
+        coordinator.cancel();
+        assert!(new_count.count() > 0);
+        assert!(other_count.count() > 0);
+        assert!(first.now_or_never().is_some());
+        assert!(second.now_or_never().is_some());
     }
 
-    /// The final chunk must be released before the stream finishes.
+    #[tokio::test]
+    async fn nlj_normal_completion_does_not_cancel_peers() -> Result<()> {
+        tokio::time::timeout(Duration::from_secs(5), async {
+            let (plan, ctx) = cancellation_test_plan()?;
+            let mut rows = 0;
+            for partition in 0..2 {
+                let batches =
+                    common::collect(plan.execute(partition, 
Arc::clone(&ctx))?).await?;
+                rows += 
batches.iter().map(RecordBatch::num_rows).sum::<usize>();
+                assert!(!plan.fallback_coordinator.is_cancelled());
+            }
+            assert_eq!(rows, 9);
+            assert_eq!(ctx.memory_pool().reserved(), 0);
+            Ok(())
+        })
+        .await
+        .expect("normal execution hung")
+    }
+
+    /// A stream that ends in an error is unfinished, so dropping it cancels 
the
+    /// peers.
     ///
-    /// `release_chunk` is what drops the coordinator's `Arc<JoinLeftData>` and
-    /// lets the next load `resize(0)` the coordinator reservation. For every
+    /// The error path reaches `Drop` without passing through `Done`. The build
+    /// side has to resolve for the poll to get as far as the failing right 
input,
+    /// so this uses a `LeftLoad::Spilled` future rather than a pending one, 
and
+    /// asserts the injected error actually surfaced before the drop -- 
otherwise
+    /// the test would silently degrade into "an unstarted stream cancels", 
which
+    /// other tests already cover.
+    #[tokio::test]
+    async fn nlj_errored_stream_drop_cancels_peers() -> Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+
+        let right_schema = build_right_table().schema();
+        let (schema, columns) =
+            build_join_schema(&spill.schema, &right_schema, &JoinType::Left);
+        let left_spill = Arc::clone(&spill);
+        let mut failing = NestedLoopJoinStream::new(
+            Arc::new(schema),
+            None,
+            JoinType::Left,
+            Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                right_schema,
+                futures::stream::once(async {
+                    exec_err!("injected right-input failure")
+                }),
+            )),
+            OnceFut::new(async move { Ok(LeftLoad::Spilled(left_spill)) }),
+            columns,
+            NestedLoopJoinMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
+            1,
+            SpillState::Pending {
+                task_context: Arc::clone(&ctx),
+                fallback_coordinator: Arc::clone(&coordinator),
+            },
+        );
+
+        // Drive it until the injected error comes out. Bounded so a fixture 
that
+        // stops reaching the right input fails instead of spinning.
+        let err = tokio::time::timeout(Duration::from_secs(5), async {
+            loop {
+                match failing.next().await {
+                    Some(Err(e)) => return e,
+                    Some(Ok(_)) => {}
+                    None => panic!("stream finished without surfacing the 
error"),
+                }
+            }
+        })
+        .await
+        .expect("fixture never reached its failing right input");
+        assert_contains!(err.to_string(), "injected right-input failure");
+        assert!(
+            !matches!(failing.state, NLJState::Done),
+            "an errored stream must not look finished"
+        );
+        assert!(!coordinator.is_cancelled());
+
+        drop(failing);
+        assert!(
+            coordinator.is_cancelled(),
+            "an unfinished stream must cancel its peers when dropped, 
including \
+             one that ended in an error"
+        );
+        Ok(())
+    }
+
+    /// A survivor parked on the global-right replay must be woken by a peer's
+    /// cancellation.
+    ///
+    /// `EmitGlobalRightUnmatched` is only reachable with `Active` spill 
state, so
+    /// the fixture installs one rather than assigning the state on top of
+    /// `Pending` -- otherwise the handler reaches its pending read through a 
path
+    /// a real execution never takes, and the test would be checking a 
`Pending`
+    /// watcher while claiming to check the replay configuration. The replay
+    /// reader is injected directly, which skips the reopen step; that is the
+    /// shortcut here, and reopening itself is covered elsewhere.
+    #[tokio::test]
+    async fn nlj_cancel_wakes_stream_parked_on_global_right_replay() -> 
Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+        let _chunk = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&ctx), Time::new())
+            .await?
+            .expect("chunk");
+
+        let right_schema = build_right_table().schema();
+        let make_stream = || {
+            let (schema, columns) =
+                build_join_schema(&spill.schema, &right_schema, 
&JoinType::Full);
+            let left_spill = Arc::clone(&spill);
+            NestedLoopJoinStream::new(
+                Arc::new(schema),
+                None,
+                JoinType::Full,
+                Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                    Arc::clone(&right_schema),
+                    futures::stream::pending::<Result<RecordBatch>>(),
+                )),
+                OnceFut::new(async move { Ok(LeftLoad::Spilled(left_spill)) }),
+                columns,
+                NestedLoopJoinMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
+                1,
+                SpillState::Pending {
+                    task_context: Arc::clone(&ctx),
+                    fallback_coordinator: Arc::clone(&coordinator),
+                },
+            )
+        };
+        let peer = make_stream();
+        let mut survivor = make_stream();
+
+        // Reach `Active` the way execution does, by letting the spilled build
+        // side resolve, then park in the replay stage.
+        let wakes = WakeCount::new();
+        let waker = futures::task::waker(Arc::clone(&wakes));
+        let mut cx = std::task::Context::from_waker(&waker);
+        assert!(survivor.poll_next_unpin(&mut cx).is_pending());
+        assert!(
+            matches!(survivor.spill_state, SpillState::Active(_)),
+            "the global-right replay only exists with Active spill state"
+        );
+
+        survivor.state = NLJState::EmitGlobalRightUnmatched;
+        survivor.left_exhausted = true;
+        // Inject the replay reader so the handler parks on it rather than
+        // reopening the spill file, which other tests cover.
+        survivor.right_data =
+            Some(Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                Arc::clone(&right_schema),
+                futures::stream::pending::<Result<RecordBatch>>(),
+            )));
+        assert!(survivor.poll_next_unpin(&mut cx).is_pending());
+        assert!(matches!(survivor.state, NLJState::EmitGlobalRightUnmatched));
+        wakes.reset();
+
+        // Poll the peer into `Active` first, so its drop exercises the real
+        // wiring: an `Active` partition disappearing cancels immediately, 
with no
+        // test standing in for the coordinated path being entered.
+        let mut peer = peer;
+        assert!(peer.poll_next_unpin(&mut cx).is_pending());
+        assert!(matches!(peer.spill_state, SpillState::Active(_)));
+        wakes.reset();
+        drop(peer);
+        assert!(coordinator.is_cancelled());
+        assert!(
+            wakes.count() > 0,
+            "cancel must wake a survivor parked on the global-right replay"
+        );
+
+        // And the wake must actually surface the cancellation, not just tick.
+        match survivor.poll_next_unpin(&mut cx) {
+            Poll::Ready(Some(Err(e))) => {
+                assert_contains!(e.to_string(), "cancelled");
+            }
+            _ => panic!("the woken survivor must report the cancellation"),
+        }
+        Ok(())
+    }
+
+    #[tokio::test]
+    async fn nlj_cancel_wakes_stream_parked_on_right_input() -> Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+        let _chunk = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&ctx), Time::new())
+            .await?
+            .expect("chunk");
+        let make_stream = || {
+            let right_schema = build_right_table().schema();
+            let (schema, columns) =
+                build_join_schema(&spill.schema, &right_schema, 
&JoinType::Left);
+            let left_spill = Arc::clone(&spill);
+            NestedLoopJoinStream::new(
+                Arc::new(schema),
+                None,
+                JoinType::Left,
+                Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                    right_schema,
+                    futures::stream::pending::<Result<RecordBatch>>(),
+                )),
+                OnceFut::new(async move { Ok(LeftLoad::Spilled(left_spill)) }),
+                columns,
+                NestedLoopJoinMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
+                1,
+                SpillState::Pending {
+                    task_context: Arc::clone(&ctx),
+                    fallback_coordinator: Arc::clone(&coordinator),
+                },
+            )
+        };
+        let peer = make_stream();
+        let mut survivor = make_stream();
+        let wakes = WakeCount::new();
+        let waker = futures::task::waker(Arc::clone(&wakes));
+        let mut cx = std::task::Context::from_waker(&waker);
+        assert!(survivor.poll_next_unpin(&mut cx).is_pending());
+        assert!(matches!(survivor.state, NLJState::FetchingRight));
+        wakes.reset();
+        // Poll the peer into `Active` first, so its drop exercises the real
+        // wiring: an `Active` partition disappearing cancels immediately, 
with no
+        // test standing in for the coordinated path being entered.
+        let mut peer = peer;
+        assert!(peer.poll_next_unpin(&mut cx).is_pending());
+        assert!(matches!(peer.spill_state, SpillState::Active(_)));
+        wakes.reset();
+        drop(peer);
+        assert!(coordinator.is_cancelled());
+        assert!(
+            wakes.count() > 0,
+            "cancel must wake the survivor waiting on right input"
+        );
+        Ok(())
+    }
+
+    /// Same terminal behaviour once a partition has already emitted output.
+    ///
+    /// A cancellation that arrives before any row is produced is the easy 
case.
+    /// Here both partitions emit first -- which is also what carries them into
+    /// the coordinated path -- and the survivor must still end in an error
+    /// followed by `None`, not resume delivering rows from a left side that is
+    /// now missing a prober.
+    ///
+    /// This does not assert that the coalescer holds a partial batch at the
+    /// moment of cancellation; nothing here forces it to. The buffered case is
+    /// covered by `nlj_cancellation_after_buffered_rows_ends_without_output`,
+    /// which establishes that premise before checking the terminal sequence.
+    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
+    async fn nlj_cancellation_after_emitting_output_ends_in_error() -> 
Result<()> {
+        let runtime = RuntimeEnvBuilder::new()
+            .with_memory_limit(50, 1.0)
+            .build_arc()?;
+        // Large enough that the coalescer batches rows rather than emitting 
each
+        // one, so this exercises a different output path than the shared
+        // `batch_size = 1` fixtures.
+        let cfg = TaskContext::default()
+            .session_config()
+            .clone()
+            .with_batch_size(8192);
+        let ctx = Arc::new(
+            TaskContext::default()
+                .with_runtime(runtime)
+                .with_session_config(cfg),
+        );
+        let right = Arc::new(RepartitionExec::try_new(
+            build_right_table_one_batch_per_row(),
+            Partitioning::RoundRobinBatch(2),
+        )?) as Arc<dyn ExecutionPlan>;
+        let plan = Arc::new(NestedLoopJoinExec::try_new(
+            build_left_table(),
+            right,
+            None,
+            &JoinType::Left,
+            None,
+        )?);
+
+        let mut peer = plan.execute(0, Arc::clone(&ctx))?;
+        let mut survivor = plan.execute(1, Arc::clone(&ctx))?;
+
+        // Both partitions produce output first, which is what carries them 
into
+        // the coordinated path -- a peer dropped while still `Pending` would 
only
+        // record the drop, since the left side might not spill at all.
+        peer.next().await.expect("peer output")?;
+        survivor.next().await.expect("survivor output")?;
+
+        drop(peer);
+
+        // Terminal behaviour must be an error, then nothing -- never a flush 
of
+        // rows produced before a partition went missing.
+        let mut saw_error = false;
+        for _ in 0..8 {
+            match survivor.next().await {
+                Some(Err(_)) => {
+                    saw_error = true;
+                    break;
+                }
+                Some(Ok(_)) => {}
+                None => break,
+            }
+        }
+        assert!(saw_error, "the survivor must report the cancellation");
+        assert!(
+            survivor.next().await.is_none(),
+            "nothing may follow the cancellation error"
+        );
+        Ok(())
+    }
+
+    #[tokio::test]
+    async fn nlj_stream_stops_producing_after_cancellation_error() -> 
Result<()> {
+        let (plan, ctx) = cancellation_test_plan()?;
+        let mut peer = plan.execute(0, Arc::clone(&ctx))?;
+        let mut survivor = plan.execute(1, Arc::clone(&ctx))?;
+        peer.next().await.expect("output")?;
+        survivor.next().await.expect("output")?;
+        drop(peer);
+        assert!(survivor.next().await.expect("cancellation").is_err());
+        let after_error = survivor.next().await;
+        assert!(
+            after_error.is_none(),
+            "cancellation did not discard buffered output"
+        );
+        Ok(())
+    }
+
+    #[tokio::test]
+    async fn nlj_cancel_before_watcher_registration_is_not_lost() -> 
Result<()> {
+        let ctx = Arc::new(TaskContext::default());
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+        {
+            let mut inner = coordinator.inner.lock();
+            inner.left_schema = Some(Arc::clone(&spill.schema));
+            inner.left_stream =
+                Some(Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                    Arc::clone(&spill.schema),
+                    futures::stream::pending::<Result<RecordBatch>>(),
+                )));
+        }
+        coordinator
+            .cancel_at_leader_claim
+            .store(1, Ordering::SeqCst);
+        let result = tokio::time::timeout(
+            Duration::from_secs(1),
+            Arc::clone(&coordinator).next_chunk(0, spill, ctx, Time::new()),
+        )
+        .await;
+        assert!(coordinator.is_cancelled());
+        assert!(
+            result.is_ok(),
+            "cancel before watcher registration was lost"
+        );
+        assert!(result.unwrap().is_err());
+        Ok(())
+    }
+
+    /// Chunk-progress notifications must not disturb a loader's cancellation
+    /// watcher.
+    ///
+    /// The re-arm gap the review found came from watching the shared `notify`,
+    /// where an unrelated wake forced a re-register and could lose a 
concurrent
+    /// cancel. The watcher now waits on a dedicated `cancel_notify` that only
+    /// `cancel` ever signals, so there is no re-arm to race: any wake is a 
real
+    /// cancellation. This drives unrelated `notify_waiters()` traffic past a
+    /// parked loader and then cancels for real.
+    #[tokio::test]
+    async fn nlj_chunk_progress_does_not_resolve_cancel_watcher() -> 
Result<()> {
+        let ctx = Arc::new(TaskContext::default());
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+        {
+            let mut inner = coordinator.inner.lock();
+            inner.left_schema = Some(Arc::clone(&spill.schema));
+            inner.left_stream =
+                Some(Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                    Arc::clone(&spill.schema),
+                    futures::stream::pending::<Result<RecordBatch>>(),
+                )));
+        }
+        let mut load = Arc::clone(&coordinator)
+            .next_chunk(0, spill, ctx, Time::new())
+            .boxed();
+        assert!(futures::poll!(load.as_mut()).is_pending());
+
+        // Unrelated progress traffic: must neither complete nor cancel the 
load.
+        for _ in 0..5 {
+            coordinator.notify.notify_waiters();
+            assert!(
+                futures::poll!(load.as_mut()).is_pending(),
+                "chunk-progress traffic must not resolve the cancellation 
watcher"
+            );
+        }
+        assert!(!coordinator.is_cancelled());
+
+        // A real cancellation must still be observed while the read is parked.
+        coordinator.cancel();
+        let result = tokio::time::timeout(Duration::from_secs(1), load)
+            .await
+            .expect("a parked loader must observe cancellation");
+        assert!(result.is_err());
+        assert!(coordinator.inner.lock().current.is_none());
+        Ok(())
+    }
+
+    /// Entering the coordinated path applies a drop recorded before it.
+    ///
+    /// This is the drop-then-coordinate direction: the deferral holds while
+    /// nobody coordinates, and `begin_coordination` converts it. The opposite
+    /// order -- a peer already `Active` when the drop happens -- runs on a 
real
+    /// plan in `nlj_pending_drop_cancels_active_peer_real_plan`, since the two
+    /// orders take different branches and only both together cover the flags.
+    #[tokio::test]
+    async fn nlj_begin_coordination_applies_prior_pending_drop() -> Result<()> 
{
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+
+        // A partition goes away while still `Pending`: recorded, not applied.
+        Arc::clone(&coordinator).record_pending_drop();
+        assert!(
+            !coordinator.is_cancelled(),
+            "a pending drop must not cancel before anything coordinates"
+        );
+
+        // A peer then enters the coordinated path and picks the drop up.
+        Arc::clone(&coordinator).begin_coordination();
+        assert!(
+            coordinator.is_cancelled(),
+            "entering the coordinated path must apply a recorded pending drop"
+        );
+        Ok(())
+    }
+
+    /// A recorded pending drop is inert if the execution never coordinates.
+    #[tokio::test]
+    async fn nlj_pending_drop_stays_inert_without_coordination() -> Result<()> 
{
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+
+        Arc::clone(&coordinator).record_pending_drop();
+        // Nobody calls `begin_coordination`, so chunks are still served.
+        let served = Arc::clone(&coordinator)
+            .next_chunk(0, spill, Arc::clone(&ctx), Time::new())
+            .await?;
+        assert!(
+            served.is_some(),
+            "a recorded pending drop must not block chunk service on its own"
+        );
+        assert!(!coordinator.is_cancelled());
+        Ok(())
+    }
+
+    /// Dropping an unfinished partition must not fail its peers when the left
+    /// side turns out to fit in memory.
+    ///
+    /// Every eligible execution passes through `SpillState::Pending` before 
the
+    /// shared load decides between `InMemory` and `Spilled`. With an ample 
pool
+    /// the load resolves to `InMemory`, there is no shared chunk counter, and 
the
+    /// right partitions stay independent -- so a dropped peer has nothing to
+    /// coordinate and must not cancel anything.
+    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
+    async fn 
nlj_drop_pending_partition_does_not_cancel_when_left_fits_in_memory()
+    -> Result<()> {
+        // Ample pool: the left side never spills, so the fallback is never 
used.
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let task_ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+        let right = Arc::new(RepartitionExec::try_new(
+            build_right_table_one_batch_per_row(),
+            Partitioning::RoundRobinBatch(2),
+        )?) as Arc<dyn ExecutionPlan>;
+        let plan = Arc::new(NestedLoopJoinExec::try_new(
+            build_left_table(),
+            right,
+            None,
+            &JoinType::Left,
+            None,
+        )?);
+
+        let abandoned = plan.execute(0, Arc::clone(&task_ctx))?;
+        let survivor = plan.execute(1, Arc::clone(&task_ctx))?;
+        drop(abandoned);
+
+        let batches = common::collect(survivor).await?;
+        assert!(
+            batches.iter().map(RecordBatch::num_rows).sum::<usize>() > 0,
+            "the survivor must still produce its rows when nothing spilled"
+        );
+        Ok(())
+    }
+
+    #[tokio::test]
+    async fn nlj_pending_drop_cancels_active_peer_real_plan() -> Result<()> {
+        tokio::time::timeout(Duration::from_secs(5), async {
+            let (plan, ctx) = cancellation_test_plan()?;
+            let pending = plan.execute(0, Arc::clone(&ctx))?;
+            let mut active = plan.execute(1, Arc::clone(&ctx))?;
+            active.next().await.expect("active output")?;
+            assert!(plan.fallback_coordinator.inner.lock().current.is_some());
+            drop(pending);
+            // The drop must have cancelled immediately: a peer is already
+            // coordinating, so this is not a deferred record.
+            assert!(
+                plan.fallback_coordinator.is_cancelled(),
+                "dropping a pending peer must cancel once a peer coordinates"
+            );
+            let result = common::collect(active).await;
+            assert!(
+                result.is_err(),
+                "already-active survivor must fail when pending peer 
disappears"
+            );
+            assert_eq!(ctx.memory_pool().reserved(), 0);
+            Ok(())
+        })
+        .await
+        .expect("survivor hung")
+    }
+
+    #[tokio::test]
+    async fn nlj_cancellation_after_buffered_rows_ends_without_output() -> 
Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+        let _chunk = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&ctx), Time::new())
+            .await?
+            .expect("chunk");
+        let right_batches =
+            common::collect(build_right_table().execute(0, 
Arc::clone(&ctx))?).await?;
+        let make_stream = |emit_input: bool| {
+            let right_schema = build_right_table().schema();
+            let (schema, columns) =
+                build_join_schema(&spill.schema, &right_schema, 
&JoinType::Left);
+            let left_spill = Arc::clone(&spill);
+            let batches = if emit_input {
+                vec![Ok(right_batches[0].slice(0, 1))]
+            } else {
+                vec![]
+            };
+            NestedLoopJoinStream::new(
+                Arc::new(schema),
+                None,
+                JoinType::Left,
+                Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                    right_schema,
+                    futures::stream::iter(batches)
+                        
.chain(futures::stream::pending::<Result<RecordBatch>>()),
+                )),
+                OnceFut::new(async move { Ok(LeftLoad::Spilled(left_spill)) }),
+                columns,
+                NestedLoopJoinMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
+                8192,
+                SpillState::Pending {
+                    task_context: Arc::clone(&ctx),
+                    fallback_coordinator: Arc::clone(&coordinator),
+                },
+            )
+        };
+        let mut peer = make_stream(false);
+        let mut survivor = make_stream(true);
+        let waker = futures::task::noop_waker();
+        let mut cx = std::task::Context::from_waker(&waker);
+        assert!(peer.poll_next_unpin(&mut cx).is_pending());
+        assert!(matches!(peer.spill_state, SpillState::Active(_)));
+        assert!(survivor.poll_next_unpin(&mut cx).is_pending());
+        let buffered = survivor.output_buffer.get_buffered_rows();
+        assert!(buffered > 0 && buffered < 8192);
+        assert!(!survivor.output_buffer.has_completed_batch());
+        drop(peer);
+        assert!(coordinator.is_cancelled());
+        assert!(matches!(
+            survivor.poll_next_unpin(&mut cx),
+            Poll::Ready(Some(Err(_)))
+        ));
+        let after = survivor.poll_next_unpin(&mut cx);
+        assert!(
+            matches!(after, Poll::Ready(None)),
+            "cancellation must terminate without an output batch"
+        );
+        Ok(())
+    }
+
+    fn cancellation_test_plan() -> Result<(Arc<NestedLoopJoinExec>, 
Arc<TaskContext>)> {
+        let runtime = RuntimeEnvBuilder::new()
+            .with_memory_limit(50, 1.0)
+            .build_arc()?;
+        let cfg = TaskContext::default()
+            .session_config()
+            .clone()
+            .with_batch_size(1);
+        let task_ctx = Arc::new(
+            TaskContext::default()
+                .with_runtime(runtime)
+                .with_session_config(cfg),
+        );
+        let right = Arc::new(RepartitionExec::try_new(
+            build_right_table_one_batch_per_row(),
+            Partitioning::RoundRobinBatch(2),
+        )?) as Arc<dyn ExecutionPlan>;
+        let plan = Arc::new(NestedLoopJoinExec::try_new(
+            build_left_table(),
+            right,
+            None,
+            &JoinType::Left,
+            None,
+        )?);
+        Ok((plan, task_ctx))
+    }
+
+    #[tokio::test]
+    async fn nlj_drop_pending_partition_does_not_retain_chunk_memory() -> 
Result<()> {
+        tokio::time::timeout(Duration::from_secs(5), async {
+            let (plan, ctx) = cancellation_test_plan()?;
+            let abandoned = plan.execute(0, Arc::clone(&ctx))?;
+            let survivor = plan.execute(1, Arc::clone(&ctx))?;
+            drop(abandoned);
+            let _ = common::collect(survivor).await;
+            assert_eq!(
+                ctx.memory_pool().reserved(),
+                0,
+                "pending partition drop must not strand the chunk"
+            );
+            Ok(())
+        })
+        .await
+        .expect("stream hung")
+    }
+
+    #[tokio::test]
+    async fn nlj_active_survivor_observes_cancel() -> Result<()> {
+        tokio::time::timeout(Duration::from_secs(5), async {
+            let (plan, ctx) = cancellation_test_plan()?;
+            let mut abandoned = plan.execute(0, Arc::clone(&ctx))?;
+            let mut survivor = plan.execute(1, Arc::clone(&ctx))?;
+            abandoned.next().await.expect("first output")?;
+            survivor.next().await.expect("first output")?;
+            assert!(
+                plan.fallback_coordinator
+                    .inner
+                    .lock()
+                    .current
+                    .as_ref()
+                    .expect("chunk")
+                    .is_last
+            );
+            drop(abandoned);
+            assert!(
+                plan.fallback_coordinator.inner.lock().cancelled,
+                "Drop must have cancelled"
+            );
+            let result = common::collect(survivor).await;
+            assert!(
+                result.is_err(),
+                "survivor holding final chunk reported success after 
cancellation"
+            );
+            Ok(())
+        })
+        .await
+        .expect("stream hung")
+    }
+
+    async fn run_paused_loader_cancellation(complete_read: bool) -> Result<()> 
{
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let pool = Arc::clone(&runtime.memory_pool);
+        let ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill = spill_left_for_test(build_left_table(), 
Arc::clone(&ctx)).await?;
+        let batches =
+            common::collect(build_left_table().execute(0, 
Arc::clone(&ctx))?).await?;
+        let batch = batches[0].clone();
+        let schema = batch.schema();
+        let (tx, rx) = tokio::sync::oneshot::channel::<()>();
+        {
+            let mut inner = coordinator.inner.lock();
+            inner.left_schema = Some(Arc::clone(&schema));
+            inner.left_stream =
+                Some(Box::pin(crate::stream::RecordBatchStreamAdapter::new(
+                    schema,
+                    futures::stream::once(async move {
+                        rx.await.expect("release the test read");
+                        Ok(batch)
+                    }),
+                )));
+        }
+        let mut load = Arc::clone(&coordinator)
+            .next_chunk(0, spill, ctx, Time::new())
+            .boxed();
+        assert!(futures::poll!(load.as_mut()).is_pending());
+        assert!(coordinator.inner.lock().loader_in_flight);
+        coordinator.cancel();
+        let _keep_sender = if complete_read {
+            tx.send(()).expect("read is waiting");
+            None
+        } else {
+            Some(tx)
+        };
+        let outcome = tokio::time::timeout(Duration::from_secs(1), load)
+            .await
+            .expect("cancelled loader still waits for its input");
+        assert!(outcome.is_err());
+        assert!(coordinator.inner.lock().current.is_none());
+        assert_eq!(pool.reserved(), 0);
+        Ok(())
+    }
+
+    #[tokio::test]
+    async fn nlj_cancelled_inflight_load_discards_its_publish() -> Result<()> {
+        run_paused_loader_cancellation(true).await
+    }
+
+    #[tokio::test]
+    async fn nlj_cancelled_inflight_load_stops_waiting() -> Result<()> {
+        run_paused_loader_cancellation(false).await
+    }
+
+    /// A partition dropped before finishing must cancel the coordinated
+    /// fallback rather than leave the survivors waiting forever.
+    ///
+    /// Chunk advancement needs every partition to report, so before this the
+    /// survivors fell through to the `notified()` wait in `next_chunk` and 
hung.
+    #[tokio::test]
+    async fn test_nlj_cancelled_partition_does_not_hang_survivors() -> 
Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let task_ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill =
+            spill_left_for_test(build_left_table(), 
Arc::clone(&task_ctx)).await?;
+
+        let (chunk_a, _) = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&task_ctx), 
Time::new())
+            .await?
+            .expect("chunk 0");
+
+        // Partition A disappears mid-probe: it never reports completion, so no
+        // emitter is elected and nothing would ever release the slot.
+        drop(chunk_a);
+        Arc::clone(&coordinator).cancel();
+
+        // The survivor must be told the execution is over, not left waiting.
+        let survivor = tokio::time::timeout(
+            Duration::from_secs(5),
+            Arc::clone(&coordinator).next_chunk(
+                1,
+                Arc::clone(&spill),
+                Arc::clone(&task_ctx),
+                Time::new(),
+            ),
+        )
+        .await;
+        assert!(
+            survivor.is_ok(),
+            "a cancelled partition must not leave the survivors hanging"
+        );
+        assert!(
+            survivor.unwrap().is_err(),
+            "the survivor should see the cancellation as an error"
+        );
+        Ok(())
+    }
+
+    /// Cancelling drops what the coordinator holds, so the chunk's reservation
+    /// is not stranded for the lifetime of a retained plan.
+    #[tokio::test]
+    async fn test_nlj_cancel_releases_coordinator_held_memory() -> Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let pool = Arc::clone(&runtime.memory_pool);
+        let task_ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let spill =
+            spill_left_for_test(build_left_table(), 
Arc::clone(&task_ctx)).await?;
+
+        let (chunk, _) = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&task_ctx), 
Time::new())
+            .await?
+            .expect("chunk 0");
+        assert!(pool.reserved() > 0);
+
+        // Streams go away without releasing the slot, then the drop guard
+        // cancels. The coordinator is still alive, as for a retained plan.
+        drop(chunk);
+        Arc::clone(&coordinator).cancel();
+        assert_eq!(
+            pool.reserved(),
+            0,
+            "cancelling must release what the coordinator was holding"
+        );
+        Ok(())
+    }
+
+    /// An already-cancelled coordinator refuses to serve chunks at all.
+    ///
+    /// This is the entry check in `next_chunk`, nothing more: cancellation
+    /// happens before any load starts. The publish-after-an-in-flight-read 
path
+    /// is covered by `nlj_cancelled_inflight_load_discards_its_publish`, 
which pauses a real
+    /// read and then completes it.
+    #[tokio::test]
+    async fn test_nlj_cancelled_coordinator_refuses_to_serve_chunks() -> 
Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let pool = Arc::clone(&runtime.memory_pool);
+        let task_ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+
+        let coordinator = Arc::new(FallbackCoordinator::new(1, true));
+        let spill =
+            spill_left_for_test(build_left_table(), 
Arc::clone(&task_ctx)).await?;
+
+        // Cancel first, then ask for a chunk: the entry check must refuse. The
+        // publish-after-an-in-flight-read path is covered by
+        // `nlj_cancelled_inflight_load_discards_its_publish`, which pauses a 
real read.
+        Arc::clone(&coordinator).cancel();
+        let result = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&task_ctx), 
Time::new())
+            .await;
+        assert!(
+            result.is_err(),
+            "a cancelled fallback must not serve chunks"
+        );
+        assert_eq!(
+            pool.reserved(),
+            0,
+            "a discarded load must not leave the reservation reinstated"
+        );
+        Ok(())
+    }
+
+    /// A coordinator nobody cancelled keeps serving chunks.
+    #[tokio::test]
+    async fn test_nlj_uncancelled_coordinator_serves_and_stays_live() -> 
Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let task_ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+        let coordinator = Arc::new(FallbackCoordinator::new(1, true));
+        let spill =
+            spill_left_for_test(build_left_table(), 
Arc::clone(&task_ctx)).await?;
+
+        // A plain sanity check that nothing marks the coordinator cancelled on
+        // its own. Real stream drops are covered by
+        // `nlj_normal_completion_does_not_cancel_peers`, which runs both 
partitions
+        // to completion and drops them.
+        let (chunk, _) = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&task_ctx), 
Time::new())
+            .await?
+            .expect("chunk 0");
+        assert!(!coordinator.inner.lock().cancelled);
+        drop(chunk);
+        Ok(())
+    }
+
+    /// The chunk's memory must be owned by the chunk's `JoinLeftData`, so it
+    /// stays accounted while *any* holder still references it.
+    ///
+    /// The emitter elected in `ProbeEnd` is the last stream to finish probing,
+    /// which is not necessarily the last to drop its `Arc<JoinLeftData>`: a
+    /// non-emitter can flush a completed output batch from
+    /// `maybe_flush_ready_batch` and return while still holding one, since 
that
+    /// return happens before `buffered_left_data = None`. Freeing the bytes 
when
+    /// the coordinator slot is released would therefore under-account memory
+    /// that is still live.
+    ///
+    /// This drives the coordinator directly so the check does not depend on
+    /// scheduling: after releasing the slot, the pool must still account for 
the
+    /// chunk while a reference is held, and drop to zero only once it is gone.
+    #[tokio::test]
+    async fn test_nlj_chunk_memory_is_owned_by_the_chunk_data() -> Result<()> {
+        let runtime = RuntimeEnvBuilder::new().build_arc()?;
+        let pool = Arc::clone(&runtime.memory_pool);
+        let task_ctx = Arc::new(TaskContext::default().with_runtime(runtime));
+
+        // One chunk, tracked bitmap, two nominal probe partitions.
+        let coordinator = Arc::new(FallbackCoordinator::new(2, true));
+        let left = build_left_table();
+        let spill = spill_left_for_test(Arc::clone(&left), 
Arc::clone(&task_ctx)).await?;
+
+        let (chunk, is_last) = Arc::clone(&coordinator)
+            .next_chunk(0, Arc::clone(&spill), Arc::clone(&task_ctx), 
Time::new())
+            .await?
+            .expect("the left side has rows, so a chunk must be produced");
+        assert!(is_last, "the fixture fits in a single chunk");
+
+        let accounted_while_held = pool.reserved();
+        assert!(
+            accounted_while_held > 0,
+            "loading a chunk must account for its batch and bitmap"
+        );
+
+        // Release the coordinator slot while still holding the chunk. This is
+        // the emitter's release: the slot is freed, but the data is alive.
+        coordinator.release_chunk(0);
+        assert_eq!(
+            pool.reserved(),
+            accounted_while_held,
+            "releasing the slot must not release memory that is still 
referenced"
+        );
+
+        // Dropping the last reference is what returns the bytes.
+        drop(chunk);
+        assert_eq!(
+            pool.reserved(),
+            0,
+            "dropping the last chunk reference must release its memory"
+        );
+        Ok(())
+    }
+
+    /// The final chunk must be released before the stream finishes.
+    ///
+    /// `release_chunk` is what drops the coordinator's `Arc<JoinLeftData>` and
+    /// lets the next load `resize(0)` the coordinator reservation. For every
     /// non-final chunk that happens on the way back to `BufferingLeft`, but 
the
     /// final chunk used to create the release future and then transition
     /// straight to `Done` (LEFT) or `EmitGlobalRightUnmatched` (FULL), neither
@@ -5827,6 +7101,13 @@ pub(crate) mod tests {
         assert_final_chunk_released(JoinType::Full).await
     }
 
+    /// Run a NLJ across 4 right partitions, collecting every output
+    /// partition CONCURRENTLY. This is required for the multi-chunk
+    /// coordinator path: a chunk is not released until all partitions
+    /// finish probing it, and a partition cannot advance to the next chunk
+    /// until the current one is released. Collecting partitions
+    /// sequentially would therefore deadlock; concurrent collection mirrors
+    /// how partitions actually run under the runtime.
     async fn multi_partition_memory_limited_join_collect_concurrent(
         left: Arc<dyn ExecutionPlan>,
         right: Arc<dyn ExecutionPlan>,


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to