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

   ### Describe the bug
   
   With sort-merge join (`datafusion.optimizer.prefer_hash_join = false`), a 
`LEFT`, `RIGHT` or `FULL` join with an extra join filter fails when the 
NULL-padded side has `NOT NULL` columns that the filter references:
   
   ```
   Arrow error: Invalid argument error: Column 'w' is declared as non-nullable 
but contains null values
   ```
   
   Hash join returns the correct result for the same queries, and so does 
sort-merge join when the columns are nullable. Whether it fails depends on how 
rows land in partitions: the example below succeeds with `target_partitions = 
1` and fails with 2 or 4.
   
   ### To Reproduce
   
   ```sql
   SET datafusion.optimizer.prefer_hash_join = false;
   SET datafusion.execution.target_partitions = 2;
   
   CREATE TABLE a (k INT NOT NULL, v INT NOT NULL) AS VALUES (1, 10), (2, 20), 
(3, 30);
   CREATE TABLE b (k INT NOT NULL, w INT NOT NULL) AS VALUES (1, 100), (3, 300);
   CREATE TABLE c (k INT NOT NULL, w INT NOT NULL) AS VALUES (1, 100), (2, 
200), (3, 300);
   
   SELECT * FROM a LEFT JOIN b ON a.k = b.k AND a.v < b.w ORDER BY a.k;
   -- Arrow error: Invalid argument error: Column 'w' is declared as 
non-nullable but contains null values
   
   SELECT * FROM a FULL JOIN b ON a.k = b.k AND a.v < b.w ORDER BY a.k;
   -- same error
   
   SELECT * FROM b RIGHT JOIN c ON b.k = c.k AND b.w <= c.w ORDER BY c.k;
   -- same error
   ```
   
   ### Expected behavior
   
   The same rows as hash join (`SET datafusion.optimizer.prefer_hash_join = 
true`), e.g. for the `LEFT JOIN`:
   
   ```
   +---+----+------+------+
   | k | v  | k    | w    |
   +---+----+------+------+
   | 1 | 10 | 1    | 100  |
   | 2 | 20 | NULL | NULL |
   | 3 | 30 | 3    | 300  |
   +---+----+------+------+
   ```
   
   and for the `RIGHT JOIN`:
   
   ```
   +------+------+---+-----+
   | k    | w    | k | w   |
   +------+------+---+-----+
   | 1    | 100  | 1 | 100 |
   | NULL | NULL | 2 | 200 |
   | 3    | 300  | 3 | 300 |
   +------+------+---+-----+
   ```
   
   ### Additional context
   
   Reproduced on current `main` (Oct 3 2026). Not bisected.
   
   From reading `sort_merge_join/materializing_stream.rs`: when a streamed row 
has no match while the buffered side still holds a batch, 
`null_join_streamed_row` appends it with `Some(scanning_batch_idx)`, so it is 
materialized in the same chunk as matched rows. The join filter is then 
evaluated on `RecordBatch::try_new(f.schema(), filter_columns)`, where the 
buffered-side columns are NULL-padded for that row but the filter schema 
declares them non-nullable, so batch validation fails. One possible fix is to 
build the filter batch with the buffered-side fields marked nullable; 
NULL-padded rows are null-joined whatever the filter returns.
   
   #21197 was the `LeftMark` variant of this error; it was closed when mark 
joins moved to the bitwise stream (#21184), so the outer-join path still has it.
   


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