fix: UnionExec now conforms each batch to the union's declared schema - #23861
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #23861 +/- ##
==========================================
+ Coverage 80.71% 80.87% +0.16%
==========================================
Files 1090 1101 +11
Lines 370339 375837 +5498
Branches 370339 375837 +5498
==========================================
+ Hits 298926 303967 +5041
- Misses 53603 53757 +154
- Partials 17810 18113 +303 ☔ 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.
|
@dariocurr |
An Opus-model review of 131796d caught three loose ends: - The SchemaConformingStream/conform_stream_schema doc comment claimed InterleaveExec::try_new re-validates child field data types the same way UnionExec::try_new does via calculate_union. It doesn't -- InterleaveExec::compute_properties builds fresh EquivalenceProperties without that check. Corrected the comment to state the real invariant (InterleaveExec inputs come from an already-validated UnionExec via the optimizer) and note that a genuine mismatch would now surface as an explicit error rather than a silently corrupt batch. - test_interleave_conforms_batch_schema and union_all_widening_cast_also_fixes_nullable iterated collected batches without asserting the batch list was non-empty, so both could pass vacuously if no rows came through. - union_nullable_spill.rs's regression test doc comment asserted UnionExec returns child streams "without schema coercion" -- no longer true after this fix. Updated it to explain the test now guards the SpillManager-level fix (apache#21292) specifically, since UnionExec itself no longer produces mismatched-nullability batches.
|
Follow-up in 4af8735, fixing a few loose ends from a closer pass:
Tests, fmt, and clippy still clean. |
There was a problem hiding this comment.
@dariocurr
Thanks for the follow-up.
The earlier feedback has been addressed:
InterleaveExecnow conforms each child stream before passing it toCombinedRecordBatchStream, and the direct regression test confirms that every emitted batch uses the interleave schema.- The SQL test module comment now focuses on the invariant and the related issue link.
- The documentation guarantee now matches the actual constructor and runtime behavior.
- Both batch loops that were previously vacuous now verify that output is non-empty.
- The spill test comment is up to date.
I did not find any new defects or have any inline comments.
Looks good to me. Approving.
|
This is a nice change -- thank you @dariocurr and @kosiew -- ideally we could coerce the schema at plan time but I don't know how to coerce nullability |
|
Following up on the schema-coercion-at-plan-time idea from this thread: opened #24094, which moves the re-stamping from inline logic in |
Follow-up to apache#23861. Per review feedback, moves the nullability re-stamping from inline logic in UnionExec/InterleaveExec::execute() into an explicit CoerceSchemaExec plan node, inserted by try_new above any child whose schema disagrees with the computed union schema. The coercion is now visible in EXPLAIN instead of hidden at runtime. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Follow-up to apache#23861 (issue apache#15394). Moves the schema re-stamping for nullability-mismatched UNION ALL / INTERLEAVE inputs out of `execute()` and into plan construction, via a new `CoerceSchemaExec` node inserted by `UnionExec::try_new`/`InterleaveExec::try_new` whenever a child's own schema disagrees with the computed union schema. The node is now visible in `EXPLAIN` output, is transparent for statistics/pushdown/proto purposes, and adds no measurable overhead versus inline re-stamping.
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.