Describe the bug
PiecewiseMergeJoin's classic-join output path materializes intermediate batches whose size is bounded only by the buffered side's total row count, not batch_size, and none of that memory is charged to a MemoryReservation.
In classic_join.rs:
let count = buffered_len - buffer_idx; — for streamed.x < buffered.y, one streamed row matching the whole remaining buffered side makes count = the full buffered row count.
UInt32Array::from_value(streamed_range.0, count) + take_record_batch(...) then replicate that single streamed row count times (the buffered columns are zero-copy slices; the streamed columns and the indices array are real O(count) allocations). Broadcasting one row via a constant-index take is also the slowest way to do it — variable-length values are re-copied count times through the random-access gather path.
- The whole O(buffered_len) batch is pushed into the
BatchCoalescer, which copies all rows again into its internal completed batches before re-splitting them to batch_size.
- The
next_completed_batch() check happens after the allocation, so it cannot bound the peak.
The operator does reserve memory for the buffered input side (in fact it double-charges it: once per incoming batch, then again for the concatenated batch + key arrays, with no shrink) — but the output/intermediate path has zero accounting. So the pool cannot see or bound the O(buffered_len × streamed_row_width) spike.
To Reproduce
Test + temporary probe in the classic_join.rs tests (manual exec construction, same style as the existing PWMJ unit tests): buffered side 200,000 rows × 3 Int32 columns, 2 streamed rows each greater than every buffered value, 8 MB pool of which ~5.6 MB is already taken by the buffered-side reservation:
PWMJ PROBE: intermediate batch rows=200000 (batch_size is typically 8192),
newly-allocated stream-side bytes=2,400,288 (+800,096 B indices) [x2 matches]
PWMJ VERDICT: query produced 400,000 rows (9,633,792 B of output) under an
8,388,608 B pool with ~2.4 MB headroom — none of the intermediate or output
allocations were accounted
count = 200,000 vs batch_size = 8,192 confirms the unbounded shape; the query allocating ~7 MB per match against 2.4 MB of pool headroom without an error confirms the accounting hole.
Repro test (probe prints the intermediate sizes)
#[tokio::test]
async fn pwmj_output_memory_not_charged_to_pool() -> Result<()> {
use datafusion_execution::runtime_env::RuntimeEnvBuilder;
const BUFFERED_ROWS: i32 = 200_000;
let vals: Vec<i32> = (0..BUFFERED_ROWS).collect();
let left = build_table(("a1", &vals), ("b1", &vals), ("c1", &vals));
// Streamed rows greater than every buffered value: each matches the
// whole buffered side (buffered.b1 < streamed.b2).
let right = build_table(
("a2", &vec![7, 8]),
("b2", &vec![BUFFERED_ROWS + 1, BUFFERED_ROWS + 2]),
("c2", &vec![70, 80]),
);
let on = (
Arc::new(Column::new_with_schema("b1", &left.schema())?) as _,
Arc::new(Column::new_with_schema("b2", &right.schema())?) as _,
);
let limit = 8 * 1024 * 1024;
let runtime = RuntimeEnvBuilder::new()
.with_memory_limit(limit, 1.0)
.build_arc()?;
let task_ctx = Arc::new(TaskContext::default().with_runtime(runtime));
let join = join(left, right, on, Operator::Lt, JoinType::Inner)?;
let stream = join.execute(0, task_ctx)?;
match common::collect(stream).await {
Ok(batches) => {
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
let out_bytes: usize = batches
.iter()
.map(crate::spill::get_record_batch_memory_size)
.sum();
assert_eq!(rows, 2 * BUFFERED_ROWS as usize);
eprintln!(
"PWMJ VERDICT CONFIRMED: query produced {rows} rows ({out_bytes} B \
of output) under an {limit} B pool of which ~5.6 MB was already \
reserved for the buffered side; the intermediate materialization \
reported by PWMJ PROBE above was never charged"
);
}
Err(e) => {
eprintln!("PWMJ VERDICT REFUTED: accounting held: {e}");
panic!("PWMJ output memory appears to be accounted: {e}");
}
}
Ok(())
}
}
Expected behavior
- Clamp
count to the coalescer's remaining capacity (BatchProcessState already carries start_buffer_idx/start_stream_idx for a mid-match resume), so intermediates are O(batch_size).
- Broadcast the streamed row with
slice(row_idx, 1) + repeat (e.g. ScalarValue::to_array_of_size) instead of a constant-index take.
- Charge the output path to a reservation (and stop double-charging the buffered input).
Additional context
Severity is tempered by reachability: PWMJ is only planned from SQL when optimizer.enable_piecewise_merge_join = true (default false), for a single inequality range predicate with no equijoin keys. Verified present on current main (the new PiecewiseMergeJoinBufferedFold reservations at tip belong to the min/max-extreme path, not the classic output path). Found during a joins audit; repro available. I plan to follow up with a fix PR.
Describe the bug
PiecewiseMergeJoin's classic-join output path materializes intermediate batches whose size is bounded only by the buffered side's total row count, notbatch_size, and none of that memory is charged to aMemoryReservation.In
classic_join.rs:let count = buffered_len - buffer_idx;— forstreamed.x < buffered.y, one streamed row matching the whole remaining buffered side makescount= the full buffered row count.UInt32Array::from_value(streamed_range.0, count)+take_record_batch(...)then replicate that single streamed rowcounttimes (the buffered columns are zero-copy slices; the streamed columns and the indices array are real O(count) allocations). Broadcasting one row via a constant-indextakeis also the slowest way to do it — variable-length values are re-copiedcounttimes through the random-access gather path.BatchCoalescer, which copies all rows again into its internal completed batches before re-splitting them tobatch_size.next_completed_batch()check happens after the allocation, so it cannot bound the peak.The operator does reserve memory for the buffered input side (in fact it double-charges it: once per incoming batch, then again for the concatenated batch + key arrays, with no shrink) — but the output/intermediate path has zero accounting. So the pool cannot see or bound the O(buffered_len × streamed_row_width) spike.
To Reproduce
Test + temporary probe in the
classic_join.rstests (manual exec construction, same style as the existing PWMJ unit tests): buffered side 200,000 rows × 3 Int32 columns, 2 streamed rows each greater than every buffered value, 8 MB pool of which ~5.6 MB is already taken by the buffered-side reservation:count = 200,000vsbatch_size = 8,192confirms the unbounded shape; the query allocating ~7 MB per match against 2.4 MB of pool headroom without an error confirms the accounting hole.Repro test (probe prints the intermediate sizes)
Expected behavior
countto the coalescer's remaining capacity (BatchProcessStatealready carriesstart_buffer_idx/start_stream_idxfor a mid-match resume), so intermediates are O(batch_size).slice(row_idx, 1)+ repeat (e.g.ScalarValue::to_array_of_size) instead of a constant-indextake.Additional context
Severity is tempered by reachability: PWMJ is only planned from SQL when
optimizer.enable_piecewise_merge_join = true(defaultfalse), for a single inequality range predicate with no equijoin keys. Verified present on currentmain(the newPiecewiseMergeJoinBufferedFoldreservations at tip belong to the min/max-extreme path, not the classic output path). Found during a joins audit; repro available. I plan to follow up with a fix PR.