From 7d15afa2643b7cceea0021d1b0ca119b4be47cd7 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 12:40:49 -0700 Subject: [PATCH 01/11] Cancel the NLJ coordinated fallback when a partition is dropped unfinished The partitions of a coordinated memory-limited fallback are not independent: chunk advancement needs every one of them to report. A partition that goes away before finishing therefore never reports, so no emitter is elected, nothing releases the coordinator slot, and the survivors fall through to the `notified()` wait in `next_chunk` and hang. Because the coordinator is owned by the exec rather than the streams, it also keeps holding the chunk it published, so its reservation stays charged for as long as the plan is alive. Reproduced directly against the coordinator: one partition takes chunk 0 and disappears, and a second partition asking for chunk 1 never returns. Three changes, which only make sense together: 1. The coordinator lock becomes a synchronous `parking_lot::Mutex`. Every critical section here was already synchronous -- the one slow operation, `load_one_chunk`, runs after the guard is dropped -- so `next_chunk` is restructured to decide under the lock, release it, then act: a `Decision` enum carries `Serve` / `Finished` / `Load` / `Wait` / `Cancelled` out of the locked block. This is what lets cleanup complete synchronously instead of depending on a future that a dropped stream takes with it. 2. A `cancelled` flag plus `cancel()`, driven from `Drop for NestedLoopJoinStream`. A stream that reached `Done` finished its work and cancels nothing; any other state cancels the fallback, dropping `current`, `carryover`, the shared left stream and the coordinator's reservation, and waking the waiters so they return an error rather than blocking. The guard exists from stream construction, so it also covers a cancellation that happens before the stream ever takes a chunk. Chunks other partitions still hold stay accounted until they release them. 3. A loader that was reading outside the lock must not publish into a coordinator that `cancel` has already cleaned up; doing so would reinstate the stream and reservation it just dropped. The publish path now discards its result and reports the cancellation instead. This accepts a deliberate semantic: cancelling one unfinished partition cancels the coordinated execution. It can no longer produce a complete result, so failing the remaining partitions is more honest than leaving them hanging or letting them report success from partial input. Making the release synchronous also removes the machinery that existed only to drive it across polls: `chunk_release_in_flight`, its three poll sites, `handle_releasing_final_chunk`, the `NLJState::ReleasingFinalChunk` state, and the `cx` argument `handle_emit_left_unmatched` no longer needs. With no future to interrupt, the "release dropped while pending" shape disappears on its own. Tests cover the cancellation timings that matter: a dropped partition not hanging the survivors, cancel releasing what the coordinator held, cancellation while a load is in flight discarding its result, and a normally-completed stream not cancelling anything. Each fails if `cancel()` is neutered -- the hang test by timing out, the memory test with 505 bytes still reserved. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 626 +++++++++++------- 1 file changed, 386 insertions(+), 240 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index fd2e0ae21c2ea..c7499b1d7373f 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}; @@ -1338,17 +1338,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, @@ -1425,6 +1414,16 @@ 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, + /// 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. @@ -1441,7 +1440,7 @@ 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, + inner: Mutex, /// Notified when a new chunk becomes available, when the left stream /// is exhausted, or when a chunk is released. notify: tokio::sync::Notify, @@ -1452,7 +1451,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, @@ -1461,6 +1460,7 @@ impl FallbackCoordinator { next_chunk_index: 0, current: None, loader_in_flight: false, + cancelled: false, }), notify: tokio::sync::Notify::new(), } @@ -1469,17 +1469,40 @@ impl FallbackCoordinator { /// 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, 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, 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(); + } + + /// 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. + fn cancel(self: &Arc) { + { + 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(); } @@ -1497,135 +1520,174 @@ impl FallbackCoordinator { // `spill_data` is already resolved: every partition receives the // same `Arc` from the shared `OnceAsync`, // 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 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, - &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; + }; + let mut reservation = match reservation { + Some(reservation) => reservation, + None => { + MemoryConsumer::new("NestedLoopJoinFallbackChunk".to_string()) + .with_can_spill(true) + .register(task_context.memory_pool()) + } + }; - match load_result { - Ok(LoadOutcome::Chunk { - data, - is_last, - carryover, - }) => { - inner.carryover = carryover; - if is_last { - inner.left_exhausted = true; + let load_result = Arc::clone(&self) + .load_one_chunk( + chunk_index, + &mut left_stream, + &mut reservation, + carryover, + Arc::clone(&left_schema), + build_time.clone(), + ) + .await; + + // 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" + ); } - 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); } } } - - // 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; } } @@ -1736,6 +1798,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, 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, + reservation: Option, + carryover: Option, + 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, @@ -1821,9 +1903,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, - /// In-flight chunk release future. Created by `EmitLeftUnmatched` - /// when the last partition for a chunk has finished its work. - chunk_release_in_flight: Option>, } impl SpillStateActive { @@ -2152,17 +2231,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); @@ -2210,6 +2279,29 @@ 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; + } + if let SpillState::Active(active) = &self.spill_state { + Arc::clone(&active.coordinator).cancel(); + } + } +} + impl NestedLoopJoinStream { #[expect(clippy::too_many_arguments)] pub(crate) fn new( @@ -2311,7 +2403,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 @@ -2376,18 +2467,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); @@ -2655,27 +2734,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>>> { // 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) @@ -2706,14 +2770,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; } @@ -2722,18 +2783,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; @@ -2759,39 +2809,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` 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>>> { - 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 @@ -4255,7 +4272,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, @@ -4263,7 +4280,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); } @@ -4293,7 +4310,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") }), @@ -5683,6 +5699,136 @@ pub(crate) mod tests { })) } + /// 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(()) + } + + /// Cancelling while a loader is mid-read must not let it publish into the + /// coordinator it just cleaned up. + /// + /// The loader takes the shared stream and reservation out of the slot, reads + /// without the lock, then puts them back. Without the cancellation check on + /// that publish, the resources `cancel` dropped would be reinstated and the + /// retention would come straight back. + #[tokio::test] + async fn test_nlj_cancel_during_load_discards_the_result() -> 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 let a load complete: the publish must be discarded. + 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 stream that finished normally must not cancel anything on drop. + #[tokio::test] + async fn test_nlj_completed_stream_drop_does_not_cancel() -> 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?; + + // Serving a chunk still works after a completed stream would have been + // dropped, i.e. nothing marked the coordinator cancelled. + 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. /// @@ -5722,7 +5868,7 @@ pub(crate) mod tests { // 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; + coordinator.release_chunk(0); assert_eq!( pool.reserved(), accounted_while_held, From e6eca12f504739e7d6ada14360fc7383128e147f Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 13:00:42 -0700 Subject: [PATCH 02/11] Close the cancellation gaps found in review Review of the previous commit found three places where cancellation was not actually observed, plus a hole in my own test coverage. Drop only handled `SpillState::Active`, so a stream cancelled while still `Pending` -- set up for the fallback but not yet holding a chunk -- reported nothing and cancelled nothing. That is exactly the "cancelled before taking a chunk" case the guard was supposed to cover, so the claim that placing it at stream construction was sufficient was wrong. Both states now cancel; only `Disabled` does not. Cancellation was only checked inside `next_chunk`, which a stream holding the final chunk never calls again. Such a survivor ran to completion and reported success built from an execution that had lost a partition. `poll_next` now checks on every iteration. `cancel()` could not interrupt a load already in flight: the loader awaits `load_one_chunk` directly, and waking `notify` does not reach a read parked on its input. The load now races the cancellation signal, and on cancellation drops its local stream and reservation, clears the leader claim and reports the cancellation rather than waiting for a read that may never finish. The publish check alone prevented reinstatement but not the stall. The test gap is the more important lesson. My four tests all called `coordinator.cancel()` directly, so neutering `cancel()` failed them while neutering the entire `Drop` body left all 68 tests passing -- they exercised the function but never the wiring that calls it. The four review tests added here drive real plans and streams, and two of them fail when `Drop` is neutered. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 244 +++++++++++++++++- 1 file changed, 232 insertions(+), 12 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index c7499b1d7373f..17ddff0a609d9 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1491,6 +1491,11 @@ impl FallbackCoordinator { /// 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. + /// True once a partition was dropped unfinished, cancelling the fallback. + fn is_cancelled(&self) -> bool { + self.inner.lock().cancelled + } + fn cancel(self: &Arc) { { let mut inner = self.inner.lock(); @@ -1628,16 +1633,52 @@ impl FallbackCoordinator { } }; - let load_result = Arc::clone(&self) - .load_one_chunk( - chunk_index, - &mut left_stream, - &mut reservation, - carryover, - Arc::clone(&left_schema), - build_time.clone(), - ) - .await; + // 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. + let cancelled = self.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(), + ); + let load_result = { + let mut load = std::pin::pin!(load); + let mut cancelled = std::pin::pin!(cancelled); + loop { + tokio::select! { + biased; + result = &mut load => break Some(result), + () = &mut cancelled => { + if self.inner.lock().cancelled { + break None; + } + // Woken for some other reason; keep reading. + cancelled.set(self.notify.notified()); + } + } + } + }; + 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; + } + self.notify.notify_waiters(); + 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 @@ -2090,6 +2131,18 @@ impl Stream for NestedLoopJoinStream { cx: &mut std::task::Context<'_>, ) -> Poll> { 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) && self.fallback_cancelled() { + self.state = NLJState::Done; + return Poll::Ready(Some(exec_err!( + "NestedLoopJoin coordinated fallback was cancelled because a \ + partition was dropped before finishing" + ))); + } + match self.state { // # NLJState transitions // --> FetchingRight @@ -2296,8 +2349,20 @@ impl Drop for NestedLoopJoinStream { if matches!(self.state, NLJState::Done) { return; } - if let SpillState::Active(active) = &self.spill_state { - Arc::clone(&active.coordinator).cancel(); + // Both states matter. `Active` means this stream was probing chunks; + // `Pending` means the fallback was set up but this stream had not taken + // one yet, and a partition that vanishes at that point still owes the + // report the others are waiting on. + let coordinator = match &self.spill_state { + SpillState::Active(active) => Some(Arc::clone(&active.coordinator)), + SpillState::Pending { + fallback_coordinator, + .. + } => Some(Arc::clone(fallback_coordinator)), + SpillState::Disabled => None, + }; + if let Some(coordinator) = coordinator { + coordinator.cancel(); } } } @@ -2414,6 +2479,23 @@ impl NestedLoopJoinStream { // ==== State handler functions ==== + /// True if a peer partition cancelled the coordinated fallback. + /// + /// Checked on every `poll_next` iteration: a stream holding the final chunk + /// never asks for another one, so `next_chunk`'s own check would not be + /// reached and it would otherwise report success built from an execution + /// that lost one of its partitions. + fn fallback_cancelled(&self) -> bool { + match &self.spill_state { + SpillState::Active(active) => active.coordinator.is_cancelled(), + SpillState::Pending { + fallback_coordinator, + .. + } => fallback_coordinator.is_cancelled(), + SpillState::Disabled => false, + } + } + /// Handle BufferingLeft state - prepare left side batches. /// /// In standard mode, uses OnceFut to load all left data at once. @@ -5699,6 +5781,144 @@ pub(crate) mod tests { })) } + fn review_cancellation_plan() -> Result<(Arc, Arc)> { + 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; + let plan = Arc::new(NestedLoopJoinExec::try_new( + build_left_table(), + right, + None, + &JoinType::Left, + None, + )?); + Ok((plan, task_ctx)) + } + + #[tokio::test] + async fn review_drop_pending_retains_chunk() -> Result<()> { + tokio::time::timeout(Duration::from_secs(5), async { + let (plan, ctx) = review_cancellation_plan()?; + let abandoned = plan.execute(0, Arc::clone(&ctx))?; + let survivor = plan.execute(1, Arc::clone(&ctx))?; + drop(abandoned); + let result = common::collect(survivor).await; + eprintln!( + "survivor succeeded={}, retained={}", + result.is_ok(), + ctx.memory_pool().reserved() + ); + assert_eq!( + ctx.memory_pool().reserved(), + 0, + "pending partition drop must not strand the chunk" + ); + Ok(()) + }) + .await + .expect("stream hung") + } + + #[tokio::test] + async fn review_active_survivor_must_observe_cancel() -> Result<()> { + tokio::time::timeout(Duration::from_secs(5), async { + let (plan, ctx) = review_cancellation_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 review_paused_loader(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 review_inflight_publish_is_discarded() -> Result<()> { + review_paused_loader(true).await + } + + #[tokio::test] + async fn review_inflight_load_observes_cancellation() -> Result<()> { + review_paused_loader(false).await + } + /// A partition dropped before finishing must cancel the coordinated /// fallback rather than leave the survivors waiting forever. /// From aac35cf52b33b06d69d0767ef50cc0edf33362d7 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 14:14:13 -0700 Subject: [PATCH 03/11] Deliver cancellation through a dedicated broadcast with registered waiters Follow-up review found two remaining hang paths, both traceable to using the chunk-progress `Notify` to carry cancellation. `Notify` broadcasts to waiters that already exist and stores no state, so a cancellation landing before a waiter registered was simply lost, and a task parked on something other than chunk progress had no waker registered at all. Cancellation now travels on its own `cancel_notify`, and every waiter registers with `Notified::enable()` before reading the `cancelled` flag. A cancellation arriving between those two steps is delivered rather than lost, which closes the gap at initial registration. Because `cancel_notify` is signalled only by `cancel()`, a wake from it always means a real cancellation: the re-arm loop that re-registered after unrelated wakes is gone, and with it the race that loop carried. For streams, checking a flag was not enough. A stream returning Pending from its right input is not waiting on the coordinator, so dropping a peer never woke it and the poll-loop check was not reached until some unrelated event happened to poll the task. Streams now hold a registered `cancellation_watcher` across polls and poll it each iteration, so the waker really is with the coordinator. Writing a replacement for the now-inapplicable re-arm test caught a third bug of my own: the loader's watcher was still constructed from the progress `Notify`, so any chunk-progress traffic resolved it and failed the load as if cancelled. The new test drives unrelated `notify_waiters()` past a parked loader and requires it to stay pending, then cancels for real. The obsolete re-arm test is dropped -- it drove a loop that no longer exists -- and the replacement needs no instrumentation. The hook for cancelling before watcher registration stays, since that gap still needs guarding. Verified: 8 review tests pass; joins 1162; memory_limit 38 including #24746's regressions; clippy clean. Neutering `Drop` fails 4 of the review tests, and weakening the stream check back to a plain flag read fails 1, so both mechanisms are covered rather than merely present. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 282 ++++++++++++++++-- 1 file changed, 249 insertions(+), 33 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index 17ddff0a609d9..bb8897de0b9b1 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1444,6 +1444,16 @@ pub(crate) struct FallbackCoordinator { /// Notified when a new chunk becomes available, when the left stream /// is exhausted, or when a chunk is released. notify: tokio::sync::Notify, + /// Notified once, permanently, when the fallback is cancelled. + /// + /// Separate from `notify` because cancellation is terminal and 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 register with + /// `Notified::enable` and then re-check `cancelled` under the lock, so a + /// cancellation landing between those two steps still wakes them. + cancel_notify: tokio::sync::Notify, + #[cfg(test)] + review_cancel_before_watch: AtomicUsize, } impl FallbackCoordinator { @@ -1463,6 +1473,9 @@ impl FallbackCoordinator { cancelled: false, }), notify: tokio::sync::Notify::new(), + cancel_notify: tokio::sync::Notify::new(), + #[cfg(test)] + review_cancel_before_watch: AtomicUsize::new(0), } } @@ -1492,10 +1505,36 @@ impl FallbackCoordinator { /// them; what this drops is the coordinator's own state, which nothing will /// come back for. /// True once a partition was dropped unfinished, cancelling the fallback. + /// + /// Production code observes cancellation through `cancellation_watcher` so a + /// waker is registered; this plain read is for assertions. + #[cfg(test)] fn is_cancelled(&self) -> bool { self.inner.lock().cancelled } + /// A future that resolves when the fallback is cancelled, already queued for + /// the broadcast. + /// + /// Callers that park on something other than chunk progress -- a stream + /// waiting on its right input, for instance -- have to hold one of these and + /// poll it alongside their own work, or a peer's cancellation never reaches + /// their waker. Registration happens here, before the caller re-checks + /// `cancelled`, so a cancellation in between is delivered rather than lost. + fn cancellation_watcher(self: &Arc) -> 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() + } + fn cancel(self: &Arc) { { let mut inner = self.inner.lock(); @@ -1509,6 +1548,7 @@ impl FallbackCoordinator { inner.reservation = None; } self.notify.notify_waiters(); + self.cancel_notify.notify_waiters(); } /// Fetch `expected_chunk_index`, becoming leader to load it from the @@ -1637,7 +1677,13 @@ impl FallbackCoordinator { // 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. - let cancelled = self.notify.notified(); + // Test-only hook: a peer cancels after the leader claim, + // before this task constructs its cancellation watcher. + #[cfg(test)] + if self.review_cancel_before_watch.swap(0, Ordering::SeqCst) == 1 { + self.cancel(); + } + let cancelled = self.cancel_notify.notified(); let load = Arc::clone(&self).load_one_chunk( chunk_index, &mut left_stream, @@ -1649,17 +1695,18 @@ impl FallbackCoordinator { let load_result = { let mut load = std::pin::pin!(load); let mut cancelled = std::pin::pin!(cancelled); - loop { + // 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 => break Some(result), - () = &mut cancelled => { - if self.inner.lock().cancelled { - break None; - } - // Woken for some other reason; keep reading. - cancelled.set(self.notify.notified()); - } + result = &mut load => Some(result), + () = &mut cancelled => None, } } }; @@ -2025,6 +2072,14 @@ pub(crate) struct NestedLoopJoinStream { // ======================================================================== /// State Tracking state: NLJState, + /// 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>, + /// Output buffer holds the join result to output. It will emit eagerly when /// the threshold is reached. output_buffer: Box, @@ -2135,12 +2190,35 @@ impl Stream for NestedLoopJoinStream { // 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) && self.fallback_cancelled() { - self.state = NLJState::Done; - return Poll::Ready(Some(exec_err!( - "NestedLoopJoin coordinated fallback was cancelled because a \ - partition was dropped before finishing" - ))); + 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.cancel_watch = None; + return Poll::Ready(Some(exec_err!( + "NestedLoopJoin coordinated fallback was cancelled because \ + a partition was dropped before finishing" + ))); + } } match self.state { @@ -2394,6 +2472,7 @@ impl NestedLoopJoinStream { current_right_batch: None, current_right_batch_matched: None, state: NLJState::BufferingLeft, + cancel_watch: None, left_probe_idx: 0, left_emit_idx: 0, left_exhausted: false, @@ -2479,23 +2558,6 @@ impl NestedLoopJoinStream { // ==== State handler functions ==== - /// True if a peer partition cancelled the coordinated fallback. - /// - /// Checked on every `poll_next` iteration: a stream holding the final chunk - /// never asks for another one, so `next_chunk`'s own check would not be - /// reached and it would otherwise report success built from an execution - /// that lost one of its partitions. - fn fallback_cancelled(&self) -> bool { - match &self.spill_state { - SpillState::Active(active) => active.coordinator.is_cancelled(), - SpillState::Pending { - fallback_coordinator, - .. - } => fallback_coordinator.is_cancelled(), - SpillState::Disabled => false, - } - } - /// Handle BufferingLeft state - prepare left side batches. /// /// In standard mode, uses OnceFut to load all left data at once. @@ -5781,6 +5843,160 @@ pub(crate) mod tests { })) } + #[tokio::test] + async fn review_v2_cancel_wakes_pending_right_input() -> Result<()> { + struct WakeCount(AtomicUsize); + impl futures::task::ArcWake for WakeCount { + fn wake_by_ref(this: &Arc) { + this.0.fetch_add(1, Ordering::SeqCst); + } + } + 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::>(), + )), + 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 = Arc::new(WakeCount(AtomicUsize::new(0))); + 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.0.store(0, Ordering::SeqCst); + drop(peer); + assert!(coordinator.is_cancelled()); + assert!( + wakes.0.load(Ordering::SeqCst) > 0, + "cancel must wake the survivor waiting on right input" + ); + Ok(()) + } + + #[tokio::test] + async fn review_v2_output_after_cancellation_error() -> Result<()> { + let (plan, ctx) = review_cancellation_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; + eprintln!("after cancellation error: {after_error:?}"); + assert!( + after_error.is_none(), + "cancellation did not discard buffered output" + ); + Ok(()) + } + + #[tokio::test] + async fn review_v2_cancel_before_watch_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::>(), + ))); + } + coordinator + .review_cancel_before_watch + .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 review_v2_progress_notifications_do_not_disturb_cancel_watch() -> 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::>(), + ))); + } + 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(()) + } + fn review_cancellation_plan() -> Result<(Arc, Arc)> { let runtime = RuntimeEnvBuilder::new() .with_memory_limit(50, 1.0) From b0bfc997a14022a00d15e2e9cc9e932485e7e61b Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 14:25:13 -0700 Subject: [PATCH 04/11] Take review round 3's additional tests and documentation corrections Round 3 found no blocking correctness issue: the reviewer walked the four notification interleavings and confirmed that enabling the `Notified` before reading `cancelled` covers each one, and that a single-signaller broadcast makes the removed re-arm loop unnecessary. What it did find was documentation and test names claiming more than the code delivers. Adds the reviewer's four tests: a watcher constructed before cancellation but first polled after it, plus one constructed after; several live watchers with waker replacement and progress traffic that must not wake them; a real `NestedLoopJoinStream` parked on build input receiving a cancellation wakeup when its peer is dropped; and both partitions running to normal completion and being dropped, which must leave the coordinator usable, return all nine rows, and free the memory while the plan stays alive. That last one closes the gap my own "completed stream" test never covered. Documentation corrections, all cases of promising more than happens: - `cancellation_watcher` said the future was already queued. Registration happens on its first poll, so callers have to poll it, not just hold it. - `cancel_notify` was described as notified "once, permanently". It is a broadcast carrying no state; `cancelled` is what persists. - `cancel`'s documentation had been split by a later insertion and was sitting above `is_cancelled`. Moved back, and it now records that `cancel` is the only signaller of `cancel_notify` -- the invariant the missing re-arm loop rests on. Two of my tests were renamed because their names overstated what they did. `test_nlj_cancel_during_load_discards_the_result` cancels *before* starting the load, so it exercises the entry check, not a publish after an in-flight read; it is now `test_nlj_cancelled_coordinator_refuses_to_serve_chunks` and points at the test that does pause a real read. `test_nlj_completed_stream_drop_does_not_cancel` never built or dropped a stream; it is now `test_nlj_uncancelled_coordinator_serves_and_stays_live` and points at the reviewer's test that drops real streams. Verified: joins 1166; memory_limit 38 including #24746's regressions; clippy clean. Neutering `Drop` fails 5 review tests and weakening the stream check to a plain flag read fails 2. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 195 +++++++++++++++--- 1 file changed, 168 insertions(+), 27 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index bb8897de0b9b1..b24f16fb71d8e 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1444,13 +1444,15 @@ pub(crate) struct FallbackCoordinator { /// Notified when a new chunk becomes available, when the left stream /// is exhausted, or when a chunk is released. notify: tokio::sync::Notify, - /// Notified once, permanently, when the fallback is cancelled. + /// Broadcast signalled when the fallback is cancelled. /// - /// Separate from `notify` because cancellation is terminal and 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 register with - /// `Notified::enable` and then re-check `cancelled` under the lock, so a - /// cancellation landing between those two steps still wakes them. + /// This carries no state of its own -- `cancelled` is what persists; this + /// only wakes tasks so they go and read it. 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, #[cfg(test)] review_cancel_before_watch: AtomicUsize, @@ -1497,30 +1499,26 @@ impl FallbackCoordinator { self.notify.notify_waiters(); } - /// 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. /// True once a partition was dropped unfinished, cancelling the fallback. /// - /// Production code observes cancellation through `cancellation_watcher` so a - /// waker is registered; this plain read is for assertions. + /// 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, already queued for - /// the broadcast. + /// 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 -- have to hold one of these and - /// poll it alongside their own work, or a peer's cancellation never reaches - /// their waker. Registration happens here, before the caller re-checks - /// `cancelled`, so a cancellation in between is delivered rather than lost. + /// 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) -> BoxFuture<'static, ()> { let coordinator = Arc::clone(self); async move { @@ -1535,6 +1533,17 @@ impl FallbackCoordinator { .boxed() } + /// 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) { { let mut inner = self.inner.lock(); @@ -5843,6 +5852,134 @@ pub(crate) mod tests { })) } + #[tokio::test] + async fn review_v3_cancel_wakes_pending_build_input() -> Result<()> { + struct WakeCount(AtomicUsize); + impl futures::task::ArcWake for WakeCount { + fn wake_by_ref(this: &Arc) { + this.0.fetch_add(1, Ordering::SeqCst); + } + } + 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); + NestedLoopJoinStream::new( + Arc::new(schema), + None, + JoinType::Left, + Box::pin(crate::stream::RecordBatchStreamAdapter::new( + right_schema, + futures::stream::pending::>(), + )), + OnceFut::new(futures::future::pending::>()), + 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 = Arc::new(WakeCount(AtomicUsize::new(0))); + 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.0.store(0, Ordering::SeqCst); + drop(peer); + assert!(coordinator.is_cancelled()); + assert!( + wakes.0.load(Ordering::SeqCst) > 0, + "cancel must wake the survivor waiting on build input" + ); + Ok(()) + } + + #[tokio::test] + async fn review_v3_cancel_before_first_poll() { + 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 review_v3_broadcast_and_waker_replacement() { + struct WakeCount(AtomicUsize); + impl futures::task::ArcWake for WakeCount { + fn wake_by_ref(this: &Arc) { + this.0.fetch_add(1, Ordering::SeqCst); + } + } + let coordinator = Arc::new(FallbackCoordinator::new(2, true)); + let mut first = coordinator.cancellation_watcher(); + let mut second = coordinator.cancellation_watcher(); + let old_count = Arc::new(WakeCount(AtomicUsize::new(0))); + let new_count = Arc::new(WakeCount(AtomicUsize::new(0))); + let other_count = Arc::new(WakeCount(AtomicUsize::new(0))); + 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() + ); + 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.0.load(Ordering::SeqCst), 0); + assert_eq!(other_count.0.load(Ordering::SeqCst), 0); + coordinator.cancel(); + coordinator.cancel(); + assert!(new_count.0.load(Ordering::SeqCst) > 0); + assert!(other_count.0.load(Ordering::SeqCst) > 0); + assert!(first.now_or_never().is_some()); + assert!(second.now_or_never().is_some()); + } + + #[tokio::test] + async fn review_v3_normal_finish_does_not_cancel() -> Result<()> { + tokio::time::timeout(Duration::from_secs(5), async { + let (plan, ctx) = review_cancellation_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::(); + assert!(!plan.fallback_coordinator.is_cancelled()); + } + assert_eq!(rows, 9); + assert_eq!(ctx.memory_pool().reserved(), 0); + Ok(()) + }) + .await + .expect("normal execution hung") + } + #[tokio::test] async fn review_v2_cancel_wakes_pending_right_input() -> Result<()> { struct WakeCount(AtomicUsize); @@ -6219,7 +6356,7 @@ pub(crate) mod tests { /// that publish, the resources `cancel` dropped would be reinstated and the /// retention would come straight back. #[tokio::test] - async fn test_nlj_cancel_during_load_discards_the_result() -> Result<()> { + 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)); @@ -6228,7 +6365,9 @@ pub(crate) mod tests { let spill = spill_left_for_test(build_left_table(), Arc::clone(&task_ctx)).await?; - // Cancel first, then let a load complete: the publish must be discarded. + // Cancel first, then ask for a chunk: the entry check must refuse. The + // publish-after-an-in-flight-read path is covered by + // `review_inflight_publish_is_discarded`, 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()) @@ -6245,17 +6384,19 @@ pub(crate) mod tests { Ok(()) } - /// A stream that finished normally must not cancel anything on drop. + /// A coordinator nobody cancelled keeps serving chunks. #[tokio::test] - async fn test_nlj_completed_stream_drop_does_not_cancel() -> Result<()> { + 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?; - // Serving a chunk still works after a completed stream would have been - // dropped, i.e. nothing marked the coordinator cancelled. + // A plain sanity check that nothing marks the coordinator cancelled on + // its own. Real stream drops are covered by + // `review_v3_normal_finish_does_not_cancel`, 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? From aed0eab5ab9adebe29eadbc0438ad42300a24314 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 14:39:39 -0700 Subject: [PATCH 05/11] Cover the last two cancellation parking spots Round 3 listed these as optional, sharing the already-reviewed watcher and Drop mechanisms. They are cheap and they close the two cases where I was reasoning rather than testing. A stream can park in `EmitGlobalRightUnmatched` rather than `BufferingLeft`: that state reopens the spilled right side and polls it. Since the watcher is polled at the top of every `poll_next` iteration regardless of state, it should be woken there too, and now that is asserted rather than assumed -- a FULL join parked on a pending replay records a wake when its peer is dropped. An unfinished stream that ends in an error also reaches `Drop` without passing through `Done`, so it cancels its peers. That follows from the `Done` check, but the error path had no test of its own; one now injects a failing right input and confirms the peers are cancelled. Both fail if `Drop` is neutered. Dropping `Pending` back out of `Drop`'s match -- the round-1 defect -- now fails 5 review tests rather than the 1 it did then. Verified: 14 review tests; joins 1168; memory_limit 38 including #24746's regressions; clippy clean. While adding these I removed a `#[tokio::test]` that belonged to `review_v3_cancel_wakes_pending_build_input` and restored it, then audited every test in the module for a missing attribute. Only the shared `review_paused_loader` helper lacks one, correctly. This is the second time an attribute has gone missing to careless editing here, and a silently unregistered test is worse than a failing one. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 122 ++++++++++++++++++ 1 file changed, 122 insertions(+) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index b24f16fb71d8e..b764cf7d34459 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -5852,6 +5852,128 @@ pub(crate) mod tests { })) } + /// A survivor parked on the global-right replay must be woken by a peer's + /// cancellation. + /// + /// `EmitGlobalRightUnmatched` reopens the spilled right side and polls it, + /// so a stream can sit there rather than in `BufferingLeft`. The watcher is + /// polled at the top of every iteration regardless of state, so this should + /// behave like the build-input case; it is the last parking spot left to + /// confirm. + #[tokio::test] + async fn review_v4_cancel_wakes_parked_global_right_replay() -> Result<()> { + struct WakeCount(AtomicUsize); + impl futures::task::ArcWake for WakeCount { + fn wake_by_ref(this: &Arc) { + this.0.fetch_add(1, Ordering::SeqCst); + } + } + 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"); + + // FULL tracks unmatched right rows, so the stream has a global-right + // replay stage to park in. + let make_stream = || { + let right_schema = build_right_table().schema(); + let (schema, columns) = + build_join_schema(&spill.schema, &right_schema, &JoinType::Full); + NestedLoopJoinStream::new( + Arc::new(schema), + None, + JoinType::Full, + Box::pin(crate::stream::RecordBatchStreamAdapter::new( + right_schema, + futures::stream::pending::>(), + )), + OnceFut::new(futures::future::pending::>()), + 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(); + + // Park it in the replay stage rather than in BufferingLeft. + survivor.state = NLJState::EmitGlobalRightUnmatched; + survivor.left_exhausted = true; + + let wakes = Arc::new(WakeCount(AtomicUsize::new(0))); + 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::EmitGlobalRightUnmatched)); + wakes.0.store(0, Ordering::SeqCst); + + drop(peer); + assert!(coordinator.is_cancelled()); + assert!( + wakes.0.load(Ordering::SeqCst) > 0, + "cancel must wake a survivor parked on the global-right replay" + ); + Ok(()) + } + + /// A stream that ends in an error is unfinished, so dropping it cancels the + /// peers. + /// + /// The error path reaches `Drop` without passing through `Done`, and the + /// coordinated execution has lost a partition either way, so the survivors + /// must be failed rather than left to report success from partial input. + #[tokio::test] + async fn review_v4_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); + // A right input that fails on its first poll. + let failing = NestedLoopJoinStream::new( + Arc::new(schema), + None, + JoinType::Left, + Box::pin(crate::stream::RecordBatchStreamAdapter::new( + Arc::clone(&right_schema), + futures::stream::once(async { + Err(datafusion_common::DataFusionError::Execution( + "injected right-input failure".to_string(), + )) + }), + )), + OnceFut::new(futures::future::pending::>()), + columns, + NestedLoopJoinMetrics::new(&ExecutionPlanMetricsSet::new(), 0), + 1, + SpillState::Pending { + task_context: Arc::clone(&ctx), + fallback_coordinator: Arc::clone(&coordinator), + }, + ); + + assert!(!coordinator.is_cancelled()); + // The stream never reaches `Done`, so dropping it counts as unfinished. + drop(failing); + assert!( + coordinator.is_cancelled(), + "an unfinished stream must cancel its peers when dropped, including \ + one that ended in an error" + ); + Ok(()) + } + #[tokio::test] async fn review_v3_cancel_wakes_pending_build_input() -> Result<()> { struct WakeCount(AtomicUsize); From ab130dc4c6a83e40bc7548ee6107b19645966da8 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 15:54:53 -0700 Subject: [PATCH 06/11] Make the two new fixtures actually reach what they claim Review round 4 found that neither fixture added in the previous commit exercised its stated scenario, and both were rewritten rather than patched. The error fixture built a stream with a failing right input and dropped it without polling, so it only repeated "an unstarted stream cancels its peers", which other tests already cover. Polling alone would not have helped: its `left_data` was a permanently pending future, so it could never leave `BufferingLeft` to reach the right input at all. It now resolves the build side through `LeftLoad::Spilled`, drives the stream under a bounded timeout until the injected error surfaces, and asserts that error and a non-`Done` state *before* dropping -- so a fixture that stops reaching the right input fails instead of quietly degrading. The replay fixture assigned `EmitGlobalRightUnmatched` on top of `SpillState::Pending`, a combination execution never produces. It reached the pending read only because `right_data` was already `Some`, bypassing the `Active`-only reopen branch, so it was really testing a `Pending` watcher while claiming to test the replay configuration. It now polls until the spilled build side puts it in `Active`, asserts that, and only then parks in the replay stage. The replay reader is still injected rather than reopened -- that shortcut is now stated in the doc comment -- and the test additionally requires the wake to surface the cancellation rather than merely counting a wakeup. Also corrects the rustdoc above `test_nlj_cancelled_coordinator_refuses_to_serve_chunks`, which still described cancelling mid-load and validating publication after the rename fixed the name and the body comment, and softens the `cancel_notify` wording: a delivered broadcast is itself sufficient, so it is not merely a nudge to go and read the flag. On test registration: `cargo test -- --list` reports all 14 review tests, which is the right check. The regex sweep I ran was the wrong tool -- the module has legitimate async helpers, so requiring an attribute on every `async fn` would be wrong -- and it is not committed. Verified: 14 review tests; joins 1168; memory_limit 38 including #24746's regressions; clippy clean. Neutering `Drop` fails both fixtures; removing the watcher poll fails the replay one only, which is the expected split since the error fixture does not depend on the watcher. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 282 ++++++++++-------- 1 file changed, 162 insertions(+), 120 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index b764cf7d34459..a49290703757a 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1446,8 +1446,10 @@ pub(crate) struct FallbackCoordinator { notify: tokio::sync::Notify, /// Broadcast signalled when the fallback is cancelled. /// - /// This carries no state of its own -- `cancelled` is what persists; this - /// only wakes tasks so they go and read it. Kept separate from `notify` + /// 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 @@ -5852,16 +5854,8 @@ pub(crate) mod tests { })) } - /// A survivor parked on the global-right replay must be woken by a peer's - /// cancellation. - /// - /// `EmitGlobalRightUnmatched` reopens the spilled right side and polls it, - /// so a stream can sit there rather than in `BufferingLeft`. The watcher is - /// polled at the top of every iteration regardless of state, so this should - /// behave like the build-input case; it is the last parking spot left to - /// confirm. #[tokio::test] - async fn review_v4_cancel_wakes_parked_global_right_replay() -> Result<()> { + async fn review_v3_cancel_wakes_pending_build_input() -> Result<()> { struct WakeCount(AtomicUsize); impl futures::task::ArcWake for WakeCount { fn wake_by_ref(this: &Arc) { @@ -5876,17 +5870,14 @@ pub(crate) mod tests { .next_chunk(0, Arc::clone(&spill), Arc::clone(&ctx), Time::new()) .await? .expect("chunk"); - - // FULL tracks unmatched right rows, so the stream has a global-right - // replay stage to park in. let make_stream = || { let right_schema = build_right_table().schema(); let (schema, columns) = - build_join_schema(&spill.schema, &right_schema, &JoinType::Full); + build_join_schema(&spill.schema, &right_schema, &JoinType::Left); NestedLoopJoinStream::new( Arc::new(schema), None, - JoinType::Full, + JoinType::Left, Box::pin(crate::stream::RecordBatchStreamAdapter::new( right_schema, futures::stream::pending::>(), @@ -5903,33 +5894,103 @@ pub(crate) mod tests { }; let peer = make_stream(); let mut survivor = make_stream(); - - // Park it in the replay stage rather than in BufferingLeft. - survivor.state = NLJState::EmitGlobalRightUnmatched; - survivor.left_exhausted = true; - let wakes = Arc::new(WakeCount(AtomicUsize::new(0))); 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::EmitGlobalRightUnmatched)); + assert!(matches!(survivor.state, NLJState::BufferingLeft)); wakes.0.store(0, Ordering::SeqCst); - drop(peer); assert!(coordinator.is_cancelled()); assert!( wakes.0.load(Ordering::SeqCst) > 0, - "cancel must wake a survivor parked on the global-right replay" + "cancel must wake the survivor waiting on build input" ); Ok(()) } + #[tokio::test] + async fn review_v3_cancel_before_first_poll() { + 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 review_v3_broadcast_and_waker_replacement() { + struct WakeCount(AtomicUsize); + impl futures::task::ArcWake for WakeCount { + fn wake_by_ref(this: &Arc) { + this.0.fetch_add(1, Ordering::SeqCst); + } + } + let coordinator = Arc::new(FallbackCoordinator::new(2, true)); + let mut first = coordinator.cancellation_watcher(); + let mut second = coordinator.cancellation_watcher(); + let old_count = Arc::new(WakeCount(AtomicUsize::new(0))); + let new_count = Arc::new(WakeCount(AtomicUsize::new(0))); + let other_count = Arc::new(WakeCount(AtomicUsize::new(0))); + 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() + ); + 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.0.load(Ordering::SeqCst), 0); + assert_eq!(other_count.0.load(Ordering::SeqCst), 0); + coordinator.cancel(); + coordinator.cancel(); + assert!(new_count.0.load(Ordering::SeqCst) > 0); + assert!(other_count.0.load(Ordering::SeqCst) > 0); + assert!(first.now_or_never().is_some()); + assert!(second.now_or_never().is_some()); + } + + #[tokio::test] + async fn review_v3_normal_finish_does_not_cancel() -> Result<()> { + tokio::time::timeout(Duration::from_secs(5), async { + let (plan, ctx) = review_cancellation_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::(); + 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. /// - /// The error path reaches `Drop` without passing through `Done`, and the - /// coordinated execution has lost a partition either way, so the survivors - /// must be failed rather than left to report success from partial input. + /// 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 review_v4_errored_stream_drop_cancels_peers() -> Result<()> { let runtime = RuntimeEnvBuilder::new().build_arc()?; @@ -5940,20 +6001,18 @@ pub(crate) mod tests { let right_schema = build_right_table().schema(); let (schema, columns) = build_join_schema(&spill.schema, &right_schema, &JoinType::Left); - // A right input that fails on its first poll. - let failing = NestedLoopJoinStream::new( + let left_spill = Arc::clone(&spill); + let mut failing = NestedLoopJoinStream::new( Arc::new(schema), None, JoinType::Left, Box::pin(crate::stream::RecordBatchStreamAdapter::new( - Arc::clone(&right_schema), + right_schema, futures::stream::once(async { - Err(datafusion_common::DataFusionError::Execution( - "injected right-input failure".to_string(), - )) + exec_err!("injected right-input failure") }), )), - OnceFut::new(futures::future::pending::>()), + OnceFut::new(async move { Ok(LeftLoad::Spilled(left_spill)) }), columns, NestedLoopJoinMetrics::new(&ExecutionPlanMetricsSet::new(), 0), 1, @@ -5963,8 +6022,26 @@ pub(crate) mod tests { }, ); + // 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()); - // The stream never reaches `Done`, so dropping it counts as unfinished. + drop(failing); assert!( coordinator.is_cancelled(), @@ -5974,8 +6051,18 @@ pub(crate) mod tests { 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 review_v3_cancel_wakes_pending_build_input() -> Result<()> { + async fn review_v4_cancel_wakes_parked_global_right_replay() -> Result<()> { struct WakeCount(AtomicUsize); impl futures::task::ArcWake for WakeCount { fn wake_by_ref(this: &Arc) { @@ -5990,19 +6077,21 @@ pub(crate) mod tests { .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 right_schema = build_right_table().schema(); let (schema, columns) = - build_join_schema(&spill.schema, &right_schema, &JoinType::Left); + build_join_schema(&spill.schema, &right_schema, &JoinType::Full); + let left_spill = Arc::clone(&spill); NestedLoopJoinStream::new( Arc::new(schema), None, - JoinType::Left, + JoinType::Full, Box::pin(crate::stream::RecordBatchStreamAdapter::new( - right_schema, + Arc::clone(&right_schema), futures::stream::pending::>(), )), - OnceFut::new(futures::future::pending::>()), + OnceFut::new(async move { Ok(LeftLoad::Spilled(left_spill)) }), columns, NestedLoopJoinMetrics::new(&ExecutionPlanMetricsSet::new(), 0), 1, @@ -6014,92 +6103,46 @@ pub(crate) mod tests { }; 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 = Arc::new(WakeCount(AtomicUsize::new(0))); 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)); + 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::>(), + ))); + assert!(survivor.poll_next_unpin(&mut cx).is_pending()); + assert!(matches!(survivor.state, NLJState::EmitGlobalRightUnmatched)); wakes.0.store(0, Ordering::SeqCst); + drop(peer); assert!(coordinator.is_cancelled()); assert!( wakes.0.load(Ordering::SeqCst) > 0, - "cancel must wake the survivor waiting on build input" + "cancel must wake a survivor parked on the global-right replay" ); - Ok(()) - } - #[tokio::test] - async fn review_v3_cancel_before_first_poll() { - 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 review_v3_broadcast_and_waker_replacement() { - struct WakeCount(AtomicUsize); - impl futures::task::ArcWake for WakeCount { - fn wake_by_ref(this: &Arc) { - this.0.fetch_add(1, Ordering::SeqCst); + // 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"), } - let coordinator = Arc::new(FallbackCoordinator::new(2, true)); - let mut first = coordinator.cancellation_watcher(); - let mut second = coordinator.cancellation_watcher(); - let old_count = Arc::new(WakeCount(AtomicUsize::new(0))); - let new_count = Arc::new(WakeCount(AtomicUsize::new(0))); - let other_count = Arc::new(WakeCount(AtomicUsize::new(0))); - 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() - ); - 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.0.load(Ordering::SeqCst), 0); - assert_eq!(other_count.0.load(Ordering::SeqCst), 0); - coordinator.cancel(); - coordinator.cancel(); - assert!(new_count.0.load(Ordering::SeqCst) > 0); - assert!(other_count.0.load(Ordering::SeqCst) > 0); - assert!(first.now_or_never().is_some()); - assert!(second.now_or_never().is_some()); - } - - #[tokio::test] - async fn review_v3_normal_finish_does_not_cancel() -> Result<()> { - tokio::time::timeout(Duration::from_secs(5), async { - let (plan, ctx) = review_cancellation_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::(); - assert!(!plan.fallback_coordinator.is_cancelled()); - } - assert_eq!(rows, 9); - assert_eq!(ctx.memory_pool().reserved(), 0); - Ok(()) - }) - .await - .expect("normal execution hung") + Ok(()) } #[tokio::test] @@ -6470,13 +6513,12 @@ pub(crate) mod tests { Ok(()) } - /// Cancelling while a loader is mid-read must not let it publish into the - /// coordinator it just cleaned up. + /// An already-cancelled coordinator refuses to serve chunks at all. /// - /// The loader takes the shared stream and reservation out of the slot, reads - /// without the lock, then puts them back. Without the cancellation check on - /// that publish, the resources `cancel` dropped would be reinstated and the - /// retention would come straight back. + /// 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 `review_inflight_publish_is_discarded`, 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()?; From fdb76fa0b3771f4d9c3732e1d5380363ab30ba9a Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 16:04:39 -0700 Subject: [PATCH 07/11] Correct comments the refactor left behind Review round 5 reported no blocking findings and cleared the change for upstream, but pointed at three comments still describing the pre-refactor code. The coordinator's inner state was documented as guarded by an async mutex; it is a `parking_lot::Mutex`, and the reason matters enough to record: cancellation and chunk release have to finish inside one `poll_next` rather than depending on a future a dropped stream would take away, and no critical section awaits. The coordinator reservation was described as holding the current chunk's memory, registered via `initiate_fallback` -- a function that no longer exists. It only holds bytes while a load runs; `load_one_chunk` then moves them into the chunk's `JoinLeftData` with `take()`, so accounting follows the data. The comment above `buffered_left_data = None` claimed the `Arc` reaches zero once the last partition lets go. The slot holds a strong reference too, so the reservation is freed only when every holder *and* the slot release. That is precisely why the slot cannot hold the chunk weakly, which two earlier attempts here established the hard way, so the comment now says it. Also corrects an overreach in my own reporting: I had described the sqllogictest suite as blocked by the pre-existing `needless_pass_by_value` failure in the sqllogictest crate. That was wrong -- clippy escalates warnings with `-D warnings`, `cargo test` does not. Running it works: `cargo test -p datafusion-sqllogictest --test sqllogictests -- nested_loop_join joins information_schema` passes 8/8 here. The lint failure is real but separate, and reproduces on clean `upstream/main` (`3266eaa91`) with `cargo clippy -p datafusion-sqllogictest --all-targets --all-features -- -D warnings`. Verified: joins 1168; sqllogictest 8/8 on the join and information_schema files; clippy clean for datafusion-physical-plan. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 26 +++++++++++++------ 1 file changed, 18 insertions(+), 8 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index a49290703757a..472cfcdd0da0f 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1386,12 +1386,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 inside a +/// single `poll_next` 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. + /// + /// It only holds bytes while a load is in progress: on success + /// `load_one_chunk` moves them into the chunk's `JoinLeftData` with + /// `take()`, so the accounting follows the data rather than this slot. + /// Lazily registered by the first leader, once a runtime context is + /// available. reservation: Option, /// The shared left spill stream from which chunks are read. Owned by /// the coordinator so only one partition reads it at a time. @@ -2912,9 +2920,11 @@ 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 -- which is why the slot cannot hold + // the chunk weakly instead. self.buffered_left_data = None; if self.is_memory_limited() { From bec82c00f8b9c2781855625b234f54a12d4b06ca Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 16:43:45 -0700 Subject: [PATCH 08/11] Say precisely why the slot's reference is load-bearing Round 6 cleared the change for upstream and corrected three pieces of wording. The important one: I had justified the slot's strong `Arc` by saying the reservation frees only after every holder and the slot release. True, but that is a consequence, not the reason. The reason is a handoff: a faster partition can drop its reference before a slower one has taken the chunk at all, and the slow one is served from the slot (`Decision::Serve`), so the slot has to keep the chunk alive across that gap. Stated that way, the comment explains why holding it weakly cannot work -- which two earlier attempts here established by failing. "Complete inside a single `poll_next`" was wrong about cancellation, which runs from `Drop` and has neither a poll nor an await; it now says "without another poll or await". "Only holds bytes while a load is in progress" was too absolute: the error path returns the reservation to the coordinator with its bytes still accounted, since `take()` happens only on success. The comment now describes the successful path without claiming the stronger invariant. Verified: joins 1168; clippy clean for datafusion-physical-plan. Comment-only, so memory_limit and sqllogictest were not rerun. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 28 +++++++++++-------- 1 file changed, 17 insertions(+), 11 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index 472cfcdd0da0f..927eb72673313 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1388,18 +1388,18 @@ struct CurrentChunk { /// Inner state of [`FallbackCoordinator`], guarded by a synchronous mutex. /// -/// Synchronous because cancellation and chunk release have to complete inside a -/// single `poll_next` 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. +/// 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 the leader borrows to bound one chunk load. /// - /// It only holds bytes while a load is in progress: on success - /// `load_one_chunk` moves them into the chunk's `JoinLeftData` with - /// `take()`, so the accounting follows the data rather than this slot. - /// Lazily registered by the first leader, once a runtime context is - /// available. + /// 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, /// The shared left spill stream from which chunks are read. Owned by /// the coordinator so only one partition reads it at a time. @@ -2923,8 +2923,14 @@ impl NestedLoopJoinStream { // `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 -- which is why the slot cannot hold - // the chunk weakly instead. + // 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() { From abf47172e48b17419a89299f1e43721cb246e328 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 17:15:19 -0700 Subject: [PATCH 09/11] Prepare the cancellation tests for submission Review-process naming and duplication that should not reach upstream. The fourteen tests carried `review_v2`/`v3`/`v4` prefixes recording which review round produced them, which means nothing to anyone reading the file later. They are renamed after the behaviour they check, under one `nlj_` prefix: what is woken, what is released, what refuses to serve. The two shared helpers are renamed the same way (`run_paused_loader_cancellation`, `cancellation_test_plan`). The coordinator carried a test-only `review_cancel_before_watch` field with no documentation at all -- a field named after a review round, sitting in production state. It is now `cancel_at_leader_claim` and says what it is for: firing a cancellation in the window between claiming a load and registering the watcher, which is the interleaving where a lost notification would strand every other partition, and which a test cannot otherwise reach. Two `eprintln!`s left over from debugging are gone; the assertions beside them already carried the meaning. The `WakeCount` waker was defined identically in four tests and is now one helper with `new`/`count`/`reset` -- deliberately just that, not a test framework. While doing the extraction a regex of mine rewrote the helper's own body into `self.reset()` calling itself, which the test run caught as a stack overflow rather than a wrong answer. Fixed, and a reminder that mechanical renames need the result read back. Verified: 44 NLJ tests, all registered per `cargo test -- --list`; joins 1168; clippy clean. Neutering `Drop` still fails 7 of them, and disabling the renamed test seam fails the test that depends on it, so the renames did not quietly detach any coverage. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 160 +++++++++--------- 1 file changed, 81 insertions(+), 79 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index 927eb72673313..5657f4fcebc9b 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1464,8 +1464,14 @@ pub(crate) struct FallbackCoordinator { /// `cancelled`, so a cancellation landing between those two steps is /// delivered rather than lost. cancel_notify: tokio::sync::Notify, + /// Test seam for cancelling at an interleaving a test cannot otherwise hit. + /// + /// Set to `1` to make the leader cancel after claiming the load but before + /// it registers its cancellation watcher -- the window where a lost + /// notification would strand every other partition. Consumed when it fires, + /// so one store arms it once. #[cfg(test)] - review_cancel_before_watch: AtomicUsize, + cancel_at_leader_claim: AtomicUsize, } impl FallbackCoordinator { @@ -1487,7 +1493,7 @@ impl FallbackCoordinator { notify: tokio::sync::Notify::new(), cancel_notify: tokio::sync::Notify::new(), #[cfg(test)] - review_cancel_before_watch: AtomicUsize::new(0), + cancel_at_leader_claim: AtomicUsize::new(0), } } @@ -1696,10 +1702,11 @@ impl FallbackCoordinator { // 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-only hook: a peer cancels after the leader claim, - // before this task constructs its cancellation watcher. + // 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.review_cancel_before_watch.swap(0, Ordering::SeqCst) == 1 { + if self.cancel_at_leader_claim.swap(0, Ordering::SeqCst) == 1 { self.cancel(); } let cancelled = self.cancel_notify.notified(); @@ -5839,6 +5846,32 @@ pub(crate) mod tests { /// 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) { + this.0.fetch_add(1, Ordering::SeqCst); + } + } + + impl WakeCount { + fn new() -> Arc { + 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( @@ -5871,13 +5904,7 @@ pub(crate) mod tests { } #[tokio::test] - async fn review_v3_cancel_wakes_pending_build_input() -> Result<()> { - struct WakeCount(AtomicUsize); - impl futures::task::ArcWake for WakeCount { - fn wake_by_ref(this: &Arc) { - this.0.fetch_add(1, Ordering::SeqCst); - } - } + async fn nlj_cancel_wakes_stream_parked_on_build_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)); @@ -5910,23 +5937,23 @@ pub(crate) mod tests { }; let peer = make_stream(); let mut survivor = make_stream(); - let wakes = Arc::new(WakeCount(AtomicUsize::new(0))); + 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.0.store(0, Ordering::SeqCst); + wakes.reset(); drop(peer); assert!(coordinator.is_cancelled()); assert!( - wakes.0.load(Ordering::SeqCst) > 0, + wakes.count() > 0, "cancel must wake the survivor waiting on build input" ); Ok(()) } #[tokio::test] - async fn review_v3_cancel_before_first_poll() { + 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(); @@ -5937,19 +5964,13 @@ pub(crate) mod tests { } #[tokio::test] - async fn review_v3_broadcast_and_waker_replacement() { - struct WakeCount(AtomicUsize); - impl futures::task::ArcWake for WakeCount { - fn wake_by_ref(this: &Arc) { - this.0.fetch_add(1, Ordering::SeqCst); - } - } + 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 = Arc::new(WakeCount(AtomicUsize::new(0))); - let new_count = Arc::new(WakeCount(AtomicUsize::new(0))); - let other_count = Arc::new(WakeCount(AtomicUsize::new(0))); + 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)); @@ -5969,20 +5990,20 @@ pub(crate) mod tests { .is_pending() ); coordinator.notify.notify_waiters(); - assert_eq!(new_count.0.load(Ordering::SeqCst), 0); - assert_eq!(other_count.0.load(Ordering::SeqCst), 0); + assert_eq!(new_count.count(), 0); + assert_eq!(other_count.count(), 0); coordinator.cancel(); coordinator.cancel(); - assert!(new_count.0.load(Ordering::SeqCst) > 0); - assert!(other_count.0.load(Ordering::SeqCst) > 0); + assert!(new_count.count() > 0); + assert!(other_count.count() > 0); assert!(first.now_or_never().is_some()); assert!(second.now_or_never().is_some()); } #[tokio::test] - async fn review_v3_normal_finish_does_not_cancel() -> Result<()> { + async fn nlj_normal_completion_leaves_coordinator_usable() -> Result<()> { tokio::time::timeout(Duration::from_secs(5), async { - let (plan, ctx) = review_cancellation_plan()?; + let (plan, ctx) = cancellation_test_plan()?; let mut rows = 0; for partition in 0..2 { let batches = @@ -6008,7 +6029,7 @@ pub(crate) mod tests { /// the test would silently degrade into "an unstarted stream cancels", which /// other tests already cover. #[tokio::test] - async fn review_v4_errored_stream_drop_cancels_peers() -> Result<()> { + 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)); @@ -6078,13 +6099,7 @@ pub(crate) mod tests { /// reader is injected directly, which skips the reopen step; that is the /// shortcut here, and reopening itself is covered elsewhere. #[tokio::test] - async fn review_v4_cancel_wakes_parked_global_right_replay() -> Result<()> { - struct WakeCount(AtomicUsize); - impl futures::task::ArcWake for WakeCount { - fn wake_by_ref(this: &Arc) { - this.0.fetch_add(1, Ordering::SeqCst); - } - } + 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)); @@ -6122,7 +6137,7 @@ pub(crate) mod tests { // Reach `Active` the way execution does, by letting the spilled build // side resolve, then park in the replay stage. - let wakes = Arc::new(WakeCount(AtomicUsize::new(0))); + 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()); @@ -6142,12 +6157,12 @@ pub(crate) mod tests { ))); assert!(survivor.poll_next_unpin(&mut cx).is_pending()); assert!(matches!(survivor.state, NLJState::EmitGlobalRightUnmatched)); - wakes.0.store(0, Ordering::SeqCst); + wakes.reset(); drop(peer); assert!(coordinator.is_cancelled()); assert!( - wakes.0.load(Ordering::SeqCst) > 0, + wakes.count() > 0, "cancel must wake a survivor parked on the global-right replay" ); @@ -6162,13 +6177,7 @@ pub(crate) mod tests { } #[tokio::test] - async fn review_v2_cancel_wakes_pending_right_input() -> Result<()> { - struct WakeCount(AtomicUsize); - impl futures::task::ArcWake for WakeCount { - fn wake_by_ref(this: &Arc) { - this.0.fetch_add(1, Ordering::SeqCst); - } - } + 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)); @@ -6202,24 +6211,24 @@ pub(crate) mod tests { }; let peer = make_stream(); let mut survivor = make_stream(); - let wakes = Arc::new(WakeCount(AtomicUsize::new(0))); + 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.0.store(0, Ordering::SeqCst); + wakes.reset(); drop(peer); assert!(coordinator.is_cancelled()); assert!( - wakes.0.load(Ordering::SeqCst) > 0, + wakes.count() > 0, "cancel must wake the survivor waiting on right input" ); Ok(()) } #[tokio::test] - async fn review_v2_output_after_cancellation_error() -> Result<()> { - let (plan, ctx) = review_cancellation_plan()?; + 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")?; @@ -6227,7 +6236,6 @@ pub(crate) mod tests { drop(peer); assert!(survivor.next().await.expect("cancellation").is_err()); let after_error = survivor.next().await; - eprintln!("after cancellation error: {after_error:?}"); assert!( after_error.is_none(), "cancellation did not discard buffered output" @@ -6236,7 +6244,7 @@ pub(crate) mod tests { } #[tokio::test] - async fn review_v2_cancel_before_watch_is_not_lost() -> Result<()> { + 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?; @@ -6250,7 +6258,7 @@ pub(crate) mod tests { ))); } coordinator - .review_cancel_before_watch + .cancel_at_leader_claim .store(1, Ordering::SeqCst); let result = tokio::time::timeout( Duration::from_secs(1), @@ -6276,8 +6284,7 @@ pub(crate) mod tests { /// cancellation. This drives unrelated `notify_waiters()` traffic past a /// parked loader and then cancels for real. #[tokio::test] - async fn review_v2_progress_notifications_do_not_disturb_cancel_watch() -> Result<()> - { + 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?; @@ -6315,7 +6322,7 @@ pub(crate) mod tests { Ok(()) } - fn review_cancellation_plan() -> Result<(Arc, Arc)> { + fn cancellation_test_plan() -> Result<(Arc, Arc)> { let runtime = RuntimeEnvBuilder::new() .with_memory_limit(50, 1.0) .build_arc()?; @@ -6343,18 +6350,13 @@ pub(crate) mod tests { } #[tokio::test] - async fn review_drop_pending_retains_chunk() -> Result<()> { + async fn nlj_drop_while_pending_releases_chunk() -> Result<()> { tokio::time::timeout(Duration::from_secs(5), async { - let (plan, ctx) = review_cancellation_plan()?; + 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 result = common::collect(survivor).await; - eprintln!( - "survivor succeeded={}, retained={}", - result.is_ok(), - ctx.memory_pool().reserved() - ); + let _ = common::collect(survivor).await; assert_eq!( ctx.memory_pool().reserved(), 0, @@ -6367,9 +6369,9 @@ pub(crate) mod tests { } #[tokio::test] - async fn review_active_survivor_must_observe_cancel() -> Result<()> { + async fn nlj_active_survivor_observes_cancel() -> Result<()> { tokio::time::timeout(Duration::from_secs(5), async { - let (plan, ctx) = review_cancellation_plan()?; + 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")?; @@ -6399,7 +6401,7 @@ pub(crate) mod tests { .expect("stream hung") } - async fn review_paused_loader(complete_read: bool) -> Result<()> { + 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)); @@ -6444,13 +6446,13 @@ pub(crate) mod tests { } #[tokio::test] - async fn review_inflight_publish_is_discarded() -> Result<()> { - review_paused_loader(true).await + async fn nlj_cancelled_inflight_load_discards_its_publish() -> Result<()> { + run_paused_loader_cancellation(true).await } #[tokio::test] - async fn review_inflight_load_observes_cancellation() -> Result<()> { - review_paused_loader(false).await + async fn nlj_cancelled_inflight_load_stops_waiting() -> Result<()> { + run_paused_loader_cancellation(false).await } /// A partition dropped before finishing must cancel the coordinated @@ -6533,7 +6535,7 @@ pub(crate) mod tests { /// /// 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 `review_inflight_publish_is_discarded`, which pauses a real + /// 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<()> { @@ -6547,7 +6549,7 @@ pub(crate) mod tests { // Cancel first, then ask for a chunk: the entry check must refuse. The // publish-after-an-in-flight-read path is covered by - // `review_inflight_publish_is_discarded`, which pauses a real read. + // `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()) @@ -6575,7 +6577,7 @@ pub(crate) mod tests { // A plain sanity check that nothing marks the coordinator cancelled on // its own. Real stream drops are covered by - // `review_v3_normal_finish_does_not_cancel`, which runs both partitions + // `nlj_normal_completion_leaves_coordinator_usable`, 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()) From 5fdcffd3542e20eb16dc6d20993f7988c2f288a2 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Sun, 6 Sep 2026 17:23:24 -0700 Subject: [PATCH 10/11] Fix a rustdoc I detached and sharpen three descriptions Round 7 found no blocking issue but caught five precision problems, one of them mine from the previous commit. Inserting the shared `WakeCount` put it directly beneath an existing rustdoc, so it captured documentation belonging to `multi_partition_memory_limited_join_collect_concurrent` -- the block explaining why that helper must collect partitions concurrently rather than sequentially. `WakeCount` now carries only its own docs and the explanation is back on the helper it describes. Same class of mistake as the recursive-body regex last commit: an insertion that looked local but moved something adjacent. `nlj_normal_completion_leaves_coordinator_usable` claimed more than it checks -- it establishes that normal completion does not cancel peers of the same execution, not that a plan or coordinator can be re-executed (it cannot; the coordinator's state is one-shot). Renamed to `nlj_normal_completion_does_not_cancel_peers`. `nlj_drop_while_pending_releases_chunk` named a chunk release that has not happened at that point: no chunk is loaded before the drop. It asserts the memory is not retained afterwards, so it is now `nlj_drop_pending_partition_does_not_retain_chunk_memory`. Removing `Pending` from `Drop`'s match still fails it, so the name change tracks the same coverage. The test seam's documentation said it reproduced an interleaving "a test cannot otherwise hit" and that a lost cancellation would "strand every other partition". Both overstated: it reproduces the interleaving deterministically, and what it strands is the loader itself, since other observers may already be returning errors. Verified: joins 1168; clippy clean. For the record, the nested_loop_join module holds 82 tests; the `nlj_` prefix selects 44 of them, so earlier reports of "44 NLJ tests" were describing the filter, not the module. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 27 ++++++++++--------- 1 file changed, 14 insertions(+), 13 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index 5657f4fcebc9b..8231ec01508cb 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1464,11 +1464,12 @@ pub(crate) struct FallbackCoordinator { /// `cancelled`, so a cancellation landing between those two steps is /// delivered rather than lost. cancel_notify: tokio::sync::Notify, - /// Test seam for cancelling at an interleaving a test cannot otherwise hit. + /// 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 -- the window where a lost - /// notification would strand every other partition. Consumed when it fires, + /// 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, @@ -5839,13 +5840,6 @@ 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 @@ -6001,7 +5995,7 @@ pub(crate) mod tests { } #[tokio::test] - async fn nlj_normal_completion_leaves_coordinator_usable() -> Result<()> { + 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; @@ -6350,7 +6344,7 @@ pub(crate) mod tests { } #[tokio::test] - async fn nlj_drop_while_pending_releases_chunk() -> Result<()> { + 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))?; @@ -6577,7 +6571,7 @@ pub(crate) mod tests { // A plain sanity check that nothing marks the coordinator cancelled on // its own. Real stream drops are covered by - // `nlj_normal_completion_leaves_coordinator_usable`, which runs both partitions + // `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()) @@ -6717,6 +6711,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, right: Arc, From 21817bb9d34c89f5c15539a7907d21996ebe82c6 Mon Sep 17 00:00:00 2001 From: Liang-Chi Hsieh Date: Tue, 8 Sep 2026 11:10:51 -0700 Subject: [PATCH 11/11] Defer a Pending partition's drop instead of cancelling on it A stream dropped while its execution was still `SpillState::Pending` cancelled the coordinator outright. `Pending` is entered by every eligible execution before the shared load decides whether the left side spills, so a join whose left side fits in memory -- one that never coordinates at all -- was failed by a partition going away. Record the drop instead and let the coordinator decide. `pending_drop` and `coordination_started` are updated and read under the same lock, so whichever of the two happens second performs the cancellation: * drop, then coordinate -- `begin_coordination` sees the recorded drop. * coordinate, then drop -- `record_pending_drop` sees coordination started and cancels immediately, because that peer is already waiting on a probe report the departing partition will never make. The earlier fix only covered the first order, since nothing remembered that coordination had begun; a plan reaching `Active` before the drop was left hanging with the chunk still reserved. An execution that resolves to `InMemory` never calls `begin_coordination`, so the recorded drop stays inert -- which is the point. Also stop a cancelled stream from emitting after its error. `handle_done` pads an empty result with one empty batch so an all-filtered join keeps its schema; a survivor cancelled before producing rows took that path and returned `Some(Ok())` after the failure. `cancelled_terminally` ends the stream at `None` instead. Two of the watcher tests now drop a peer that has really reached `Active` rather than calling `cancel()` by hand, so they exercise the `Drop` wiring. `nlj_cancel_wakes_stream_parked_on_build_input` keeps its explicit `begin_coordination()` -- its `left_data` never resolves, so no peer there can coordinate -- and asserts the drop stays inert first. Comments and two test names that the earlier mechanism left stale are corrected in passing: the removed `pending_drops`/`apply_pending_drops` are no longer referenced, the `Drop` comment no longer claims the `Active` peer cancels itself, and neither test name now promises a premise it does not set up. Co-authored-by: Claude Code --- .../src/joins/nested_loop_join.rs | 399 +++++++++++++++++- 1 file changed, 387 insertions(+), 12 deletions(-) diff --git a/datafusion/physical-plan/src/joins/nested_loop_join.rs b/datafusion/physical-plan/src/joins/nested_loop_join.rs index 8231ec01508cb..ab0e5d23f53c1 100644 --- a/datafusion/physical-plan/src/joins/nested_loop_join.rs +++ b/datafusion/physical-plan/src/joins/nested_loop_join.rs @@ -1422,6 +1422,26 @@ 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. /// @@ -1489,6 +1509,8 @@ 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(), @@ -1550,6 +1572,55 @@ impl FallbackCoordinator { .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) { + 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) { + 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. /// @@ -2099,6 +2170,14 @@ 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 @@ -2240,6 +2319,7 @@ impl Stream for NestedLoopJoinStream { }; 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 \ @@ -2454,20 +2534,24 @@ impl Drop for NestedLoopJoinStream { if matches!(self.state, NLJState::Done) { return; } - // Both states matter. `Active` means this stream was probing chunks; - // `Pending` means the fallback was set up but this stream had not taken - // one yet, and a partition that vanishes at that point still owes the - // report the others are waiting on. - let coordinator = match &self.spill_state { - SpillState::Active(active) => Some(Arc::clone(&active.coordinator)), + // `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, .. - } => Some(Arc::clone(fallback_coordinator)), - SpillState::Disabled => None, - }; - if let Some(coordinator) = coordinator { - coordinator.cancel(); + } => Arc::clone(fallback_coordinator).record_pending_drop(), + SpillState::Disabled => {} } } } @@ -2499,6 +2583,7 @@ 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, @@ -2612,7 +2697,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)))), } } @@ -3095,6 +3188,12 @@ impl NestedLoopJoinStream { /// Handle Done state - final state processing fn handle_done(&mut self) -> Poll>> { + // 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; @@ -5938,6 +6037,17 @@ pub(crate) mod tests { 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!( + !coordinator.is_cancelled(), + "a pending drop must not cancel before anything coordinates" + ); + // 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, @@ -6153,6 +6263,13 @@ pub(crate) mod tests { 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!( @@ -6211,6 +6328,13 @@ pub(crate) mod tests { 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!( @@ -6220,6 +6344,79 @@ pub(crate) mod tests { 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; + 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()?; @@ -6316,6 +6513,184 @@ pub(crate) mod tests { 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; + 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::() > 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::>()), + )), + 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, Arc)> { let runtime = RuntimeEnvBuilder::new() .with_memory_limit(50, 1.0)