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]
