viirya opened a new issue, #24745:
URL: https://github.com/apache/datafusion/issues/24745

   ### Describe the bug
   
   A memory-limited `NestedLoopJoinExec` silently drops unmatched rows when it 
spills. The query succeeds and returns *fewer* rows than the same query with an 
ample memory pool — there is no error.
   
   This is pre-existing and not a regression from #24675: that PR only touches 
`nested_loop_join.rs`'s spill gate, and its gate does not cover this path.
   
   ### To Reproduce
   
   Two tables registered from in-memory batches (`l`: 200 rows, `r`: 90 rows, 
join keys drawn from a small domain so most rows match), `target_partitions = 
1`, `batch_size = 16`, comparing an unlimited pool against a 64-byte pool:
   
   ```sql
   -- LEFT ANTI, via NOT EXISTS
   SELECT l.k, l.v FROM l WHERE NOT EXISTS (SELECT 1 FROM r WHERE l.v > r.w);
   -- ample: 2 rows, memory-limited: 0 rows
   
   SELECT l.k, r.w FROM l LEFT JOIN r ON l.v > r.w;
   -- ample: 9336 rows, memory-limited: 9334 rows
   
   SELECT l.k, r.w FROM l FULL JOIN r ON l.v > r.w;
   -- ample: 9339 rows, memory-limited: 9337 rows
   ```
   
   The lost rows are the unmatched ones (e.g. `6|NULL`, `8|NULL` for the `LEFT 
JOIN`). `INNER` and explicit `RIGHT JOIN` are unaffected. The join spills in 
the failing runs (`spill_count > 0`).
   
   Note the failing shapes reach the operator as *swapped* join types (`Right`, 
`RightAnti`, `RightSemi`, or `Full`), so the rows needing unmatched emission 
sit on the operator's **probe** side.
   
   ### Root cause
   
   In memory-limited mode the left side is processed in chunks and the right 
side is replayed per chunk. Per-right-batch match bitmaps are therefore merged 
into `SpillStateActive::global_right_bitmaps` and their emission is 
deliberately deferred to `NLJState::EmitGlobalRightUnmatched`, which is only 
reachable from `handle_emit_left_unmatched`'s `left_exhausted == true` branch.
   
   But `handle_buffering_left` contains this early exit:
   
   ```rust
   if active.pending_batches.is_empty() {
       // No data at all — go directly to Done
       self.left_exhausted = true;
       self.state = NLJState::Done;
       return ControlFlow::Continue(());
   }
   ```
   
   Despite the comment, this also fires on the load *after the final left 
chunk*, when the left side is exhausted and nothing remains to buffer. Going 
straight to `Done` discards the accumulated bitmaps, so any probe row that no 
chunk matched is never emitted.
   
   Instrumenting the failing `LEFT JOIN` case shows 39 bitmap merges but 
**zero** emissions, and the completion path only ever running with 
`left_exhausted = false`.
   
   ### Expected behavior
   
   Memory-limited execution returns the same rows as unlimited execution — or 
fails with `ResourcesExhausted`. It must not return fewer rows silently.
   


-- 
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