This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/datafusion.git
The following commit(s) were added to refs/heads/main by this push:
new d3c47b4a80 fix: drain PiecewiseMergeJoin output before state
transitions (#24298)
d3c47b4a80 is described below
commit d3c47b4a80baa3ca941ca0fc2cc80f96687c5a27
Author: Burak Εen <[email protected]>
AuthorDate: Sat Aug 15 14:18:41 2026 +0000
fix: drain PiecewiseMergeJoin output before state transitions (#24298)
## Which issue does this PR close?
- Closes #24297
## Rationale for this change
`ClassicPWMJStream` transitions phases while its output `BatchCoalescer`
still holds completed batches. The Left/Full unmatched pass falls back
into its producer when the drained queue is an exact multiple of
`batch_size` β emitting duplicate rows forever β and queued unmatched
Right/Full batches are dropped at terminal transitions. A zero-row
placeholder batch also escapes as stream output on empty results.
## What changes are included in this PR?
- Add `BatchProcessState::next_drained_batch()` (`finish_buffered_batch`
+ `next_completed_batch`); `None` guarantees the coalescer is empty
- `process_stream_batch` and `process_unmatched_buffered_batch` drain
one batch per poll and transition only after a confirmed-empty pop, so a
finished producer can never re-run and no queued batch is lost
- Scan completion returns `StatefulStreamResult::Continue` instead of
the empty placeholder batch; empty joins now yield no batches
## Are these changes tested?
Yes. The existing unmatched Left/Right tests now run with `batch_size =
1` (the Left one bounded with `take(6)` so the old duplicate loop fails
as a snapshot mismatch instead of hanging) β both fail on main and pass
here. The empty-join tests assert no batches are emitted.
## Are there any user-facing changes?
No (PWMJ is experimental and off by default).
---
.../src/joins/piecewise_merge_join/classic_join.rs | 123 ++++++++-------------
1 file changed, 49 insertions(+), 74 deletions(-)
diff --git
a/datafusion/physical-plan/src/joins/piecewise_merge_join/classic_join.rs
b/datafusion/physical-plan/src/joins/piecewise_merge_join/classic_join.rs
index 1c89927087..bc69ae5140 100644
--- a/datafusion/physical-plan/src/joins/piecewise_merge_join/classic_join.rs
+++ b/datafusion/physical-plan/src/joins/piecewise_merge_join/classic_join.rs
@@ -277,6 +277,18 @@ impl ClassicPWMJStream {
return Ok(StatefulStreamResult::Ready(Some(batch)));
}
+ // A finished scan can leave several completed batches queued; emit
+ // them one per poll and transition only once the queue is empty, so
+ // no output is lost.
+ if !self.batch_process_state.continue_process {
+ if let Some(batch) =
self.batch_process_state.next_drained_batch()? {
+ return Ok(StatefulStreamResult::Ready(Some(batch)));
+ }
+
+ self.state = PiecewiseMergeJoinStreamState::FetchStreamBatch;
+ return Ok(StatefulStreamResult::Continue);
+ }
+
// Produce more work
let batch = resolve_classic_join(
buffered_side,
@@ -289,25 +301,8 @@ impl ClassicPWMJStream {
)?;
if !self.batch_process_state.continue_process {
- // We finished scanning this stream batch.
- self.batch_process_state
- .output_batches
- .finish_buffered_batch()?;
- if let Some(b) = self
- .batch_process_state
- .output_batches
- .next_completed_batch()
- {
- self.state = PiecewiseMergeJoinStreamState::FetchStreamBatch;
- return Ok(StatefulStreamResult::Ready(Some(b)));
- }
-
- // Nothing pending; hand back whatever `resolve` returned (often
empty) and move on.
- if self.batch_process_state.output_batches.is_empty() {
- self.state = PiecewiseMergeJoinStreamState::FetchStreamBatch;
-
- return Ok(StatefulStreamResult::Ready(Some(batch)));
- }
+ // Scan finished; re-enter through the drain guard above.
+ return Ok(StatefulStreamResult::Continue);
}
Ok(StatefulStreamResult::Ready(Some(batch)))
@@ -324,25 +319,13 @@ impl ClassicPWMJStream {
}
if !self.batch_process_state.continue_process {
- if let Some(batch) = self
- .batch_process_state
- .output_batches
- .next_completed_batch()
- {
+ if let Some(batch) =
self.batch_process_state.next_drained_batch()? {
return Ok(StatefulStreamResult::Ready(Some(batch)));
}
- self.batch_process_state
- .output_batches
- .finish_buffered_batch()?;
- if let Some(batch) = self
- .batch_process_state
- .output_batches
- .next_completed_batch()
- {
- self.state = PiecewiseMergeJoinStreamState::Completed;
- return Ok(StatefulStreamResult::Ready(Some(batch)));
- }
+ // Fully drained; finish instead of re-running the pass.
+ self.state = PiecewiseMergeJoinStreamState::Completed;
+ return Ok(StatefulStreamResult::Continue);
}
let buffered_data =
@@ -372,29 +355,8 @@ impl ClassicPWMJStream {
self.batch_process_state.output_batches.push_batch(batch)?;
self.batch_process_state.continue_process = false;
- if let Some(batch) = self
- .batch_process_state
- .output_batches
- .next_completed_batch()
- {
- return Ok(StatefulStreamResult::Ready(Some(batch)));
- }
-
- self.batch_process_state
- .output_batches
- .finish_buffered_batch()?;
- if let Some(batch) = self
- .batch_process_state
- .output_batches
- .next_completed_batch()
- {
- self.state = PiecewiseMergeJoinStreamState::Completed;
- return Ok(StatefulStreamResult::Ready(Some(batch)));
- }
-
- self.state = PiecewiseMergeJoinStreamState::Completed;
- self.batch_process_state.reset();
- Ok(StatefulStreamResult::Ready(None))
+ // Re-enter through the drain guard above.
+ Ok(StatefulStreamResult::Continue)
}
}
@@ -438,6 +400,13 @@ impl BatchProcessState {
self.continue_process = true;
self.processed_null_count = false;
}
+
+ // `None` guarantees the coalescer holds no pending rows, so the caller
+ // may safely transition state without losing output.
+ fn next_drained_batch(&mut self) -> Result<Option<RecordBatch>> {
+ self.output_batches.finish_buffered_batch()?;
+ Ok(self.output_batches.next_completed_batch())
+ }
}
impl Stream for ClassicPWMJStream {
@@ -674,7 +643,9 @@ mod tests {
use arrow_schema::{DataType, Field};
use datafusion_common::test_util::batches_to_string;
use datafusion_execution::TaskContext;
+ use datafusion_execution::config::SessionConfig;
use datafusion_physical_expr::{PhysicalExpr, expressions::Column};
+ use futures::TryStreamExt;
use insta::assert_snapshot;
use std::sync::Arc;
@@ -964,12 +935,8 @@ mod tests {
);
let (_, batches) =
join_collect(left, right, on, Operator::LtEq,
JoinType::Inner).await?;
- assert_snapshot!(batches_to_string(&batches), @r"
- +----+----+----+----+----+----+
- | a1 | b1 | c1 | a2 | b1 | c2 |
- +----+----+----+----+----+----+
- +----+----+----+----+----+----+
- ");
+ // An empty join result produces no batches at all, not an empty batch.
+ assert!(batches.is_empty());
Ok(())
}
@@ -1298,8 +1265,16 @@ mod tests {
Arc::new(Column::new_with_schema("b1", &right.schema())?) as _,
);
- let (_, batches) =
- join_collect(left, right, on, Operator::LtEq,
JoinType::Left).await?;
+ let task_ctx = Arc::new(
+ TaskContext::default()
+ .with_session_config(SessionConfig::new().with_batch_size(1)),
+ );
+ // Bound collection so the old duplicate loop becomes a snapshot
mismatch.
+ let batches = join(left, right, on, Operator::LtEq, JoinType::Left)?
+ .execute(0, task_ctx)?
+ .take(6)
+ .try_collect::<Vec<_>>()
+ .await?;
assert_snapshot!(batches_to_string(&batches), @r"
+----+----+----+----+----+----+
@@ -1344,8 +1319,12 @@ mod tests {
Arc::new(Column::new_with_schema("b1", &right.schema())?) as _,
);
- let (_, batches) =
- join_collect(left, right, on, Operator::GtEq,
JoinType::Right).await?;
+ let task_ctx = Arc::new(
+ TaskContext::default()
+ .with_session_config(SessionConfig::new().with_batch_size(1)),
+ );
+ let join = join(left, right, on, Operator::GtEq, JoinType::Right)?;
+ let batches = common::collect(join.execute(0, task_ctx)?).await?;
assert_snapshot!(batches_to_string(&batches), @r"
+----+----+----+----+----+----+
@@ -1408,12 +1387,8 @@ mod tests {
let (_, batches) =
join_collect(left, right, on, Operator::Gt,
JoinType::Inner).await?;
- assert_snapshot!(batches_to_string(&batches), @r"
- +----+----+----+----+----+----+
- | a1 | b1 | c1 | a2 | b1 | c2 |
- +----+----+----+----+----+----+
- +----+----+----+----+----+----+
- ");
+ // An empty join result produces no batches at all, not an empty batch.
+ assert!(batches.is_empty());
Ok(())
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]