fix: UnionExec now conforms each batch to the union's declared schema - #23861
fix: UnionExec now conforms each batch to the union's declared schema#23861dariocurr wants to merge 2 commits into
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #23861 +/- ##
==========================================
+ Coverage 80.71% 80.85% +0.14%
==========================================
Files 1090 1101 +11
Lines 370339 375516 +5177
Branches 370339 375516 +5177
==========================================
+ Hits 298926 303640 +4714
- Misses 53603 53775 +172
- Partials 17810 18101 +291 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
UNION ALL between an input with a NOT NULL column and one where the same column is nullable produced a valid, correctly-typed logical plan (the analyzer already OR's nullability across legs in coerce_union_schema), but UnionExec::execute() handed out each child's RecordBatches completely unchanged. The leg whose column was already NOT NULL (and so needed no CAST) kept emitting batches with a NOT NULL field, contradicting the union's own declared (nullable) schema. Any consumer that checks schema equality across batches from the same stream -- e.g. pyarrow.Table.from_batches via the Arrow C Stream FFI used by the datafusion Python bindings -- then rejects the stream with `ArrowInvalid: Schema at index N was different`, even though every individual SELECT runs fine on its own and DataFusion's own execution never errors. UnionExec::execute() now re-stamps each child's batches with the union's own schema whenever they disagree. This is always safe: the union's schema can only be more permissive than any single input's (nullability is OR'd, never narrowed), and only the Field::nullable flag changes -- the underlying array data and data type are untouched. Closes apache#15394.
c122156 to
bed030d
Compare
|
Is there any news on that? |
kosiew
left a comment
There was a problem hiding this comment.
@dariocurr
Thanks for working on this schema consistency issue.
The UnionExec fix looks useful, but the same mismatch can still occur after the optimizer replaces it with InterleaveExec. I have left one blocking comment for that path and one small test documentation suggestion.
| input_stream_vec.push(input.execute(partition, Arc::clone(&context))?); | ||
| } else { | ||
| // Do not find a partition to execute | ||
| break; |
There was a problem hiding this comment.
I think the same schema mismatch can still surface when the optimizer replaces a UnionExec with an InterleaveExec.
InterleaveExec constructs its declared schema using the same union_schema helper, but execute() passes the child streams directly to CombinedRecordBatchStream. Its poll_next implementation then returns each child RecordBatch unchanged.
Since ensure_distribution can rewrite a UnionExec to an InterleaveExec when the children are interleavable (datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs:1493), a SQL-visible UNION ALL plan may still emit batches whose schemas do not match the declared interleave or union schema.
Could you apply the same schema-conforming wrapper to each InterleaveExec child stream, or update CombinedRecordBatchStream so yielded batches are re-stamped with its declared schema?
There was a problem hiding this comment.
Good catch, thanks. Pushed 131796d: InterleaveExec::execute() now wraps each child stream through the same conform_stream_schema/SchemaConformingStream helper (extracted from the UnionExec path) before handing them to CombinedRecordBatchStream, so batches are re-stamped with the interleave's declared schema whenever nullability disagrees. Added test_interleave_conforms_batch_schema (in union.rs) covering this directly.
| // specific language governing permissions and limitations | ||
| // under the License. | ||
|
|
||
| //! Regression tests for `UNION ALL` between inputs where the same column is |
There was a problem hiding this comment.
The regression coverage is helpful. The module-level comment currently repeats much of the PR and issue description, though.
Could we trim it down to the invariant being tested and the issue link? The assertions already make the downstream PyArrow failure mode clear.
There was a problem hiding this comment.
Trimmed in 131796d — the module doc-comment is now just the invariant being tested plus the issue link.
UnionExec::execute() already re-stamped each child's batches with the union's declared schema whenever nullability disagreed, but the same mismatch was still reachable through InterleaveExec: ensure_distribution can rewrite a UnionExec into an InterleaveExec when the children are interleavable, and InterleaveExec::execute() handed CombinedRecordBatchStream the child streams unchanged. Extract the existing UnionExec fix into a small conform_stream_schema helper and apply it to each InterleaveExec child stream before combining. Also trim the union_nullable.rs module doc-comment down to the invariant and issue link, per review feedback. Addresses review comments from kosiew on PR apache#23861.
Which issue does this PR close?
nullable: trueif any of the branches produces fields withnullable: trueon the same position #15394.Rationale for this change
UNION ALLbetween an input whose column isNOT NULLand an input where the same column is nullable produces a valid, correctly-typed logical plan — the analyzer already OR's nullability across legs incoerce_union_schema(datafusion/optimizer/src/analyzer/type_coercion.rs), so the union's declared schema correctly reports the field as nullable.The bug is at execution time:
UnionExec::execute()hands out each child'sRecordBatches completely unchanged. A leg whose column was alreadyNOT NULL(and therefore needed noCASTfrom the analyzer) keeps emitting batches with aNOT NULLfield, 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 notablypyarrow.Table.from_batchesvia the Arrow C Stream FFI used by thedatafusionPython bindings — rejects the stream withArrowInvalid: Schema at index N was different, even though every individualSELECTruns fine on its own.Minimal reproducible example (Python)
The same root cause is why
#16627had to make the sqllogictestconvert_batcheshelper 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.What changes are included in this PR?
datafusion/physical-plan/src/union.rs:UnionExec::execute()now compares each child stream's schema againstUnionExec's own declared schema, and if they disagree, wraps the child stream in a small newSchemaConformingStreamthat 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 existingcoerce_union_schemadocs), and only theField::nullablemetadata 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 stayNOT NULL), and a case where one leg also needs a realCAST(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 plainUnionExec, which is what's reported in #23862 / #15394 and reproduces the Python-binding failure above.Are these changes tested?
Yes — added
datafusion/core/tests/sql/union_nullable.rswith 4 new tests. I verified each one fails with a clear schema-mismatch assertion onmain(i.e. before this fix) and passes with it applied. Also ran the fulldatafusion-physical-plananddatafusion-optimizerunit suites,union.slt/union_by_name.sltsqllogictests, andcargo fmt/clippy(--no-deps, since an unrelated pre-existing dead-code lint indatafusion-physical-exprfails-D warningsonmaineven without this change).Are there any user-facing changes?
UNION ALLresults now consistently report the analyzer's declared nullability on every batch, regardless of which leg produced it. No public API changes.