zhuqi-lucas opened a new issue, #25270:
URL: https://github.com/apache/datafusion/issues/25270

   ### 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.
   
   <details>
   <summary>Repro test (probe prints the intermediate sizes)</summary>
   
   ```rust
   #[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(())
       }
   }
   ```
   
   </details>
   
   ### Expected behavior
   
   1. 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).
   2. Broadcast the streamed row with `slice(row_idx, 1)` + repeat (e.g. 
`ScalarValue::to_array_of_size`) instead of a constant-index `take`.
   3. 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.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to