dariocurr opened a new pull request, #23861:
URL: https://github.com/apache/datafusion/pull/23861
## Which issue does this PR close?
- Closes #15394.
## Rationale for this change
`UNION ALL` between an input whose column is `NOT NULL` and an input where
the same column is nullable produces a valid, correctly-typed logical
plan — the analyzer already OR's nullability across legs in
`coerce_union_schema` (`datafusion/optimizer/src/analyzer/type_coercion.rs`),
so the union's *declared* schema correctly says the field is nullable.
The bug is at execution time: `UnionExec::execute()` handed out each
child's `RecordBatch`es completely unchanged. A leg whose column was
already `NOT NULL` (and therefore needed no `CAST` from the analyzer) kept
emitting batches with a `NOT NULL` field, contradicting the union's own
declared (nullable) schema. DataFusion's own execution tolerates this
silently, but any consumer that checks schema equality across batches
from the same stream — most notably `pyarrow.Table.from_batches` via the
Arrow C Stream FFI used by the `datafusion` Python bindings — rejects the
stream with `ArrowInvalid: Schema at index N was different`, even though
every individual `SELECT` runs fine on its own.
### Minimal reproducible example (Python)
```python
import pyarrow as pa
from datafusion import SessionContext
ctx = SessionContext()
ctx.register_record_batch(
"table_a",
pa.record_batch(
{"id": [1, 2], "status": ["ok", "ok"]},
schema=pa.schema([("id", pa.int64()), ("status", pa.string())]), #
NOT NULL
),
)
ctx.register_record_batch(
"table_b",
pa.record_batch(
{"id": [3, 4], "status": ["done", None]},
schema=pa.schema([("id", pa.int64()), ("status", pa.string(),
True)]), # nullable
),
)
df = ctx.sql(
"SELECT id, status FROM table_a UNION ALL SELECT id, status FROM table_b"
)
print(df.schema()) # status: string, nullable -- correct
df.to_pandas() # raises pyarrow.lib.ArrowInvalid:
# Schema at index 1 was different:
# id: int64 not null / status: string not null
# vs
# id: int64 not null / status: string
```
The same root cause is why `#16627` had to make the sqllogictest
`convert_batches` helper tolerant of this exact mismatch instead of
failing, and why `#15603` (stale, closed for inactivity) attempted a
similar fix at the physical-execution layer but didn't land.
I also want to flag related report I dug into while root-causing this:
[Rationale-driven diagnostic writeup with additional repro
data](https://github.com/apache/datafusion/issues/15394#issuecomment-placeholder)
— happy to attach more detail if useful for review.
## What changes are included in this PR?
- `datafusion/physical-plan/src/union.rs`: `UnionExec::execute()` now
compares each child stream's schema against `UnionExec`'s own declared
schema, and if they disagree, wraps the child stream in a small new
`SchemaConformingStream` that re-stamps every batch with the union's
schema before yielding it. This is always safe: the union's schema can
only be *more* permissive than any single input's (nullability is
combined with logical OR, never narrowed — see the existing
`coerce_union_schema` docs), and only the `Field::nullable` metadata
changes; the underlying array data and data type are untouched.
- `datafusion/core/tests/sql/union_nullable.rs` (new): regression tests
covering same-type nullable/non-nullable mismatches in both leg orders,
the "both legs NOT NULL" case (schema should stay `NOT NULL`), and a case
where one leg also needs a real `CAST` (`Int32` -> `Int64`) in addition
to the nullability fix.
`InterleaveExec` (used for sorted unions) may have an analogous issue, but
I kept this PR scoped to plain `UnionExec`, which is what's reported in
#15394 and reproduces the Python-binding failure above.
## Are these changes tested?
Yes — added `datafusion/core/tests/sql/union_nullable.rs` with 4 new
tests. I verified each one fails with a clear schema-mismatch assertion on
`main` (i.e. before this fix) and passes with it applied. Also ran the
full `datafusion-physical-plan` and `datafusion-optimizer` unit suites,
`union.slt`/`union_by_name.slt` sqllogictests, and `cargo fmt`/`clippy`
(`--no-deps`, since an unrelated pre-existing dead-code lint in
`datafusion-physical-expr` fails `-D warnings` on `main` even without this
change).
## Are there any user-facing changes?
`UNION ALL` results now consistently report the analyzer's declared
nullability on every batch, regardless of which leg produced it. No
public API changes.
--
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]