Skip to content

fix: cast native IF branches to a common type like CASE WHEN - #6458

Merged
andygrove merged 7 commits into
apache:mainfrom
andygrove:fix/native-if-common-type-6334
Oct 2, 2026
Merged

andygrove merged 7 commits into
apache:mainfrom
andygrove:fix/native-if-common-type-6334

Conversation

@andygrove

@andygrove andygrove commented Sep 30, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6334.

Found by the 1.1.0 regression audit (#6399) and tracked in #6402.

Rationale for this change

Spark only requires the two branches of IF to have the same type up to nullability, so it adds no cast when a struct field, map value or array element can be NULL in one branch and not in the other. Comet serialized each branch with its own Arrow type, and the native IfExpr declared the THEN branch's type as its output.

  • On a batch where no row took the THEN branch, IfExpr returned the ELSE array unchanged, and the projection failed with column types must match schema types.
  • On a batch that mixed the two, DataFusion's CaseExpr casts the ELSE rows to the THEN branch's type. That failed when the THEN branch was the less nullable one, with Cannot cast nullable struct field 'x' to non-nullable field, or Found unmasked nulls for non-nullable StructArray field "value" for a map.

CASE WHEN doesn't have this problem, because create_case_expr in the planner casts the branches to a common type first.

This became a 1.1.0 regression with #5452, which made folded map literals native. A folded map('z', 0) has valueContainsNull = false and a MAP<STRING, INT> column has valueContainsNull = true, so IF(c, m, map('z', 0)) and IF(m IS NULL, map('z', 0), m) fail on 1.1.0-rc1, while 1.0.0 ran them in Spark.

What changes are included in this PR?

The native planner now builds IF the way it builds CASE WHEN, through a new create_if_expr next to create_case_when. It takes the common type of the two branches from DataFusion's type_union_coercion, which is what get_coerce_type_for_case_expression folds over the CASE WHEN branches, and wraps a branch whose Arrow type differs from it in a Comet Cast, through the same coerce_branch helper that create_case_when now uses. A branch that already has the common type is left alone, so an IF whose branches agree, which is the usual case, is planned as before.

  • The THEN branch goes first, so the common type takes its field names from the THEN branch, as Spark's If.dataType does. That also covers struct field names that differ only in case, which Spark treats as the same type. Nullability is combined the same way in either order.
  • The cast uses the same options as the CASE WHEN casts, including the UTC timezone from fix: relabel CASE and COALESCE timestamp branches instead of panicking #6347.
  • The branches share a Spark type, so the casts don't change any values. They make a nested field nullable, or relabel a field name or a timestamp's timezone. The native Cast already handles this for structs (field by field), lists (rebuilt with the target element field) and maps (the relabel path in cast_map_to_map). CASE WHEN relies on the same casts.

I did this in the native planner rather than in CometIf, because that's where CASE WHEN does it and because it compares the Arrow types the branches actually produce. A JVM-side cast to Spark's If.dataType would miss a branch whose native type differs from its Spark type.

#6350 has landed, and this branch merges main. IF can't simply call create_case_when, because that takes the common type starting from the ELSE branch, which would give IF the ELSE branch's field names. So create_if_expr keeps the THEN-first order and shares only the cast.

How are these changes tested?

Every new query runs with checkSparkAnswerAndOperator, and each table is written by three INSERTs. Each INSERT writes its own files and a batch never spans files, so every query sees batches in which every row takes the THEN branch, batches in which every row takes the ELSE branch, and batches that mix the two.

  • if_nested_nullability.sql covers the constructor path, since the SQL harness disables ConstantFolding. It runs the issue's struct and map queries in both branch orders, IF(q, m, map('z', 0)) and IF(m IS NULL, map('z', 0), m) over a map column, a struct column against a struct constructor, and a struct and an array inside a map value. It also runs IF(q, array(i), array(0)) and the CASE WHEN form, which already passed. Two more queries cover struct field names that differ only in case: the plain IF, and to_json of it in both branch orders under expect_native(if), because the row comparison ignores struct field names. The file opts in to native to_json (spark.comet.expression.StructsToJson.allowIncompatible=true), since the dispatched to_json would evaluate its whole argument, the IF included, in the JVM.
  • A new CometMapExpressionSuite test covers the folded map literal path with ConstantFolding on. It runs the audit's two queries and IF(c, map('z', 0), m).
  • Unit tests for create_if_expr in case_when.rs. if_reconciles_timestamp_timezone_labels covers branches whose timestamp timezone labels differ, which no SQL query can reach now that fix: label native timestamp_seconds results as UTC timestamps #6341 and fix: keep the input's timezone label on native date_trunc results #6345 fixed the producers of mismatched labels. if_reconciles_struct_field_nullability_and_names covers nullability in both directions and a field name that differs only in case. Each evaluates a mixed batch, an all-THEN batch and an all-ELSE batch, and checks the result's type every time.

Without the fix, all 12 IF queries whose branches differ in nullability fail (9 in the SQL file and 3 in the Scala test), and so do the two field-name queries. The two queries that already passed still pass. With the fix, all of them pass. With the casts disabled, both unit tests fail.

These pass:

  • Spark 4.1 (the default profile), 625 tests: ./mvnw test -Dtest=none -Dsuites="org.apache.comet.expressions.conditional.CometIfSuite,org.apache.comet.expressions.conditional.CometCaseWhenSuite,org.apache.comet.expressions.conditional.CometCoalesceSuite,org.apache.comet.CometMapExpressionSuite,org.apache.comet.CometSqlFileTestSuite"
  • Spark 3.4, 3.5 and 4.0, 17 tests each: CometIfSuite, the new CometMapExpressionSuite test and the expressions/conditional/ SQL files.
  • cargo test -p datafusion-comet --lib planner and cargo clippy --all-targets --workspace -- -D warnings

After merging main and addressing review, these pass as well:

  • Spark 4.1 and Spark 3.5: the expressions/conditional/ SQL files, which include perf: optimize native CASE WHEN and IF (up to 12x faster) #6350's case_when_*.sql, and the new CometMapExpressionSuite test, 20 tests.
  • cargo test -p datafusion-comet-spark-expr --lib conditional_funcs (18 tests), cargo test -p datafusion-comet --lib planner (49 tests) and cargo clippy --all-targets --workspace -- -D warnings.

Spark adds no cast to IF when its branches differ only in the nullability
of a nested field, and the native IfExpr returns a branch's array unchanged
when every row of a batch takes that branch, so the projection failed with
"column types must match schema types". Cast a branch whose Arrow type
differs from the common type, as the planner already does for CASE WHEN.
@andygrove
andygrove marked this pull request as ready for review September 30, 2026 14:26
@github-actions github-actions Bot added the bug Something isn't working label Sep 30, 2026
@andygrove andygrove mentioned this pull request Sep 30, 2026
20 of 21 tasks
@andygrove andygrove added the backport-1.1 Candidate for backporting to 1.1 release branch label Sep 30, 2026

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the clear write-up. The THEN-first common type matches Spark's If.dataType, and the new queries look like they cover both the all-ELSE and mixed-batch failures. I haven't run anything, so these are expectations from reading the code.

Comment thread native/core/src/execution/planner.rs Outdated
// names, as Spark's `If` does.
let true_type = true_expr.data_type(&input_schema)?;
let false_type = false_expr.data_type(&input_schema)?;
let common_type = type_union_coercion(&true_type, &false_type);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: would it read more simply to handle the None case up front, like let Some(common_type) = type_union_coercion(&true_type, &false_type) else { return Ok(Arc::new(IfExpr::new(if_expr, true_expr, false_expr))) };? The closure then needs no inner match or shadowed common_type, and it has the same shape as create_case_when in #6350. Would it also make sense to move this into a small create_if_expr, so a unit test like case_reconciles_timestamp_timezone_labels can cover the timestamp label path? I expect that is the only practical way to reach it now that #6341 and #6345 fixed the producers of mismatched labels.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in 171e50e. create_if_expr now sits next to create_case_when in case_when.rs, with the let-else shape, and the If arm in the planner is a single call. The cast is shared through a small coerce_branch helper. I kept IF on its own coercion rather than calling create_case_when, because that one starts from the ELSE branch and would give IF the ELSE branch's field names. if_reconciles_timestamp_timezone_labels covers the label path, and if_reconciles_struct_field_nullability_and_names covers nullability both ways and a name that differs only in case. Both evaluate mixed, all-THEN and all-ELSE batches and check the result type each time, and both fail with the casts disabled.


-- a struct column and a struct constructor
query
SELECT IF(q, s, named_struct('x', 0)) FROM test_if_nested

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would it be worth one more query where the struct field names differ only in case, for example SELECT IF(q, named_struct('x', i), named_struct('X', i)) FROM test_if_nested? Spark treats those as the same type in case-insensitive mode, so I expect it to hit a cast that only relabels the field name, which none of the current queries do. I haven't run it, but I expect it to fail without this change too.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added in 3268b69, and you were right that it failed without the change, on both batch shapes: column types must match schema types on an all-ELSE batch, and DataFusion's struct cast refusing the mismatched names on a mixed one. The row comparison ignores field names, so there's also a to_json of it in both branch orders under expect_native(if), with native to_json opted in for the file, which checks that the result takes the THEN branch's name.

Native CASE WHEN has the same problem on main since #6350: it takes the common type starting from the ELSE branch, so its result carries the ELSE branch's names. That's outside this PR, so I filed #6482.

…branch cast

Move the IF branch coercion out of the planner into `create_if_expr`, next to
`create_case_when`. It handles the no-common-type case up front, and the cast
of a differing branch is now one helper that both use. IF still takes the
common type from `type_union_coercion` with the THEN branch first, so struct
field names come from the THEN branch, as in Spark's `If.dataType`.
`create_case_when` folds from the ELSE branch, so IF can't reuse it as is.

Add unit tests that plan an IF whose timestamp branches have different
timezone labels, and one whose struct field differs in nullability or in the
case of its name. They check the result type over a batch that mixes the
branches, one where every row takes THEN, and one where every row takes ELSE.
…n case

In case-insensitive mode Spark treats `named_struct('x', i)` and
`named_struct('X', i)` as one type, adds no cast to IF, and names the result's
field after the THEN branch. Without the planner's cast, this IF failed
natively on a batch where every row took ELSE ("column types must match schema
types"), and on a batch that mixed the branches, because DataFusion's struct
cast matches fields by name.

A row comparison ignores struct field names, so the file now also runs
to_json natively over the IF, in both branch orders, which fails if the result
takes the ELSE branch's name.
Every TimestampType value in a native plan is labelled UTC, so don't enshrine an Etc/UTC result
for a THEN branch that already broke that.
@parthchandra

Copy link
Copy Markdown
Contributor

Notes (Spark's IfTypeCoercion already widens int/long/decimal/string before Comet sees the plan, so only nested-nullability, field-name, and tz-label differences survive, which is exactly what this handles):

  • native/core/src/execution/planner.rs:823 — this uses SparkCastOptions::new(EvalMode::Legacy, "UTC", false), which matches create_case_expr only after fix: relabel CASE and COALESCE timestamp branches instead of panicking #6347 lands. On current main create_case_expr still uses new_without_timezone, so if this merges first, IF relabels a timestamp branch to UTC while CASE WHEN relabels to empty during that window. "UTC" is the correct choice - can you confirm the intended merge order with fix: relabel CASE and COALESCE timestamp branches instead of panicking #6347?

  • spark/src/test/resources/sql-tests/expressions/conditional/if_nested_nullability.sql:53 — every struct here is a single field struct<x:int>, so the common type only ever differs from one branch. Please add a two-field struct where a different field is nullable in each branch, e.g. IF(q, named_struct('x', i, 'y', 0), named_struct('x', 0, 'y', i)), so both branches get cast.

  • if_nested_nullability.sql (add cases) — the nesting tested is map-of-struct and map-of-array. The complex-type bar also wants array-of-struct and struct-with-a-map-or-array-field. Please add a branch pair over array<struct<x:int>> and a struct whose field is itself a map or array, exercising the recursive struct_coercion/list_coercion paths the fix relies on.

  • native/core/src/execution/planner.rs:818 — the "every TimestampType is labelled UTC" path has no test; there's no IF query with a timestamp branch anywhere. Please add an IF over a timestamp column with a non-UTC session timezone, mirroring what fix: relabel CASE and COALESCE timestamp branches instead of panicking #6347 added for CASE/COALESCE. If it's genuinely unreachable from a native plan, a sentence saying so would be clearer than the comment.

@andygrove

Copy link
Copy Markdown
Member Author

@parthchandra Added in 4101f81: the two-field struct with complementary nullability so both branches need casts, array-of-struct, and structs containing map/array fields. The nested cases cover both branch orders, with expect_native(if) and the existing all-THEN, all-ELSE and mixed batches. The expanded if_nested_nullability fixture passes locally on Spark 4.1 and 3.5 with a freshly built native library. CI is still running.

#6347 merged as fab751c and is already included in this branch, so that merge-order prerequisite is satisfied.

The mismatched timestamp-label path is unreachable through SQL after #6341/#6345: native TimestampType values carry UTC labels even in a non-UTC session. The existing if_reconciles_timestamp_timezone_labels Rust test constructs mismatched labels directly and checks mixed, all-THEN and all-ELSE batches. The PR description explains this limitation.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Native IF branches with different nested nullability or field labels could return arrays inconsistent with the declared output type.
  • Design approach: create_if_expr computes a common Arrow type and shares coerce_branch with CASE WHEN.
  • Correctness / compatibility analysis: Spark merges struct fields positionally. DataFusion sometimes matches them by name, introducing the reproducible P2 below. Relevant Spark sources were compared across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0.
  • Key design decisions: THEN-first naming and preserving branches already matching the common type are appropriate. Those unchanged branches gain no runtime cast.
  • Implementation sketch: The planner delegates IF construction to the new helper. SQL, Scala and Rust tests cover nested nullability, field labels and different batch shapes. Sharing the existing cast helper keeps the change small.
  • Behavioral changes worth calling out: Ordinary nullability reconciliation passes the regression coverage, but case-variant field permutations can now change numeric types and results.
  • Suggested improvements: Merge struct fields recursively by position and add the regression described below.

Reviewed all five files in the full diff from b38a60b37770df6efb4c66faac751160e92990bf to 4101f81d3e8975522efc3b95691e5ef95252f155. The PR remains non-draft. Routed skills: review-comet-pr and review-comet-expression-pr. Existing discussion requests are addressed, with no substantiated unresolved existing P1/P2 blocker.

Exact-head CI: 24 successful checks, 15 skipped, no failures. Run 36786846155 includes 1,894 passing Rust tests and 1,595 passing Spark 4.1 expression tests, including the new fixtures. Spark SQL suites and macOS were skipped.

Local validation: all 18 conditional Rust tests passed. A disposable exact-head native reproduction confirmed the finding against the unchanged base constructor path. Spark 3.5.9 confirmed expected schema and results. Full JVM integration and Spark SQL suites were not run locally. Disposable project files were removed and the checkout is clean.

// The coercion that `get_coerce_type_for_case_expression` folds over the branches of a CASE
// WHEN, starting from the ELSE branch. Here the THEN branch goes first, so the common type
// takes its struct field names, as Spark's `If.dataType` does.
let Some(common_type) = type_union_coercion(&true_type, &false_type) else {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Merge struct fields positionally before using DataFusion coercion. With default case-insensitive resolution, nullable i INT and d DOUBLE, Spark accepts IF(b, named_struct('x', i, 'X', CAST(5.5 AS DOUBLE)), named_struct('X', 0, 'x', d)) and returns STRUCT<x:INT,X:DOUBLE>. Because the branches differ in nullability, type_union_coercion reaches DataFusion's name-based struct merge and pairs the first integer field with the second double field. The new casts consequently change x to DOUBLE even on an all-THEN batch that worked before. With native to_json enabled, the result changes from {"x":7,"X":5.5} to {"x":7.0,"X":5.5}. Use positional recursive struct reconciliation while retaining THEN names and combining nullability.

Evidence: An exact-head Rust reproduction constructed the branches with CreateNamedStruct over b=[true,true], i=[7,NULL], d=[9.5,NULL]. The base planner path, IfExpr::new, returned Int32 values [7,NULL] and JSON {"x":7,"X":5.5}. create_if_expr returned Float64 values [7.0,NULL] and JSON {"x":7.0,"X":5.5}. Spark 3.5.9 independently returned the integer schema and original JSON. The case-distinct names pass Scala's duplicate-name check. Native JSON observability requires spark.comet.expression.StructsToJson.allowIncompatible=true. Reproduction source is retained at /tmp/review6458-repro.rs, with output in /tmp/review6458-repro.log and Spark results in /tmp/review6458-spark-reference.log.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in ee1b213 and updated the branch-1.1 backport (#6491) in 6a83c1b. Struct fields now reconcile recursively by position, retaining THEN names and combining nullability through structs, arrays and maps. Scalar coercion is unchanged.

Confirmed the exact JSON difference and a nullable-ELSE panic on both previous heads. Added native schema/value assertions and SQL coverage for the final JSON result, both branch orders, all-THEN/all-ELSE/mixed batches, nulls and nested containers. The expanded 21-query fixture passes on Spark 3.5 and 4.1 for both patches; the Rust conditional suites (19 source / 6 backport) and workspace Clippy also pass. New-head CI is pending.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Native IF could fail when branches differed in nested nullability or struct-field names because their Arrow types remained inconsistent.
  • Design approach: create_if_expr recursively reconciles structs, arrays and maps, preserving THEN field names and combining nullability. It shares coerce_branch with CASE WHEN.
  • Correctness / compatibility analysis: Compared Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Positional reconciliation matches Spark’s rules. The previously reported P2 is fixed: the native reproduction now preserves INT and produces {"x":7,"X":5.5}.
  • Key design decisions: Branches already matching the common type avoid runtime casts. Reusing existing casts keeps the implementation focused. No reproducible performance regression or unnecessary abstraction was identified.
  • Implementation sketch: The planner delegates IF construction to the new helper. Rust, SQL and Scala tests cover nested types, field names, timestamp labels, folded literals and different batch shapes.
  • Behavioral changes worth calling out: Previously failing branches now return a consistent Arrow type. Existing discussion concerns are addressed, with no substantiated unresolved P1/P2 blocker.
  • Suggested improvements: None meeting the reporting threshold. No introduced P1/P2 issues found within this review.

Reviewed all five files in the full diff from b38a60b37770df6efb4c66faac751160e92990bf to ee1b213a2f1677510d20d86a7d562ba64e169d46. The PR remains non-draft. Routed skills: review-comet-pr and review-comet-expression-pr.

Exact-head CI: 24 successful checks, 15 skipped, no failures or pending checks. Run 36854803282 passed 1,934 Rust tests and 1,616 Spark 4.1 expression tests, including the new SQL fixture and folded-map regression. These jobs tested GitHub’s merge commit 94bf585, containing the requested head.

Local validation: all 19 conditional Rust tests and the disposable JSON regression reproduction passed from the pinned checkout. Full JVM integration, Spark SQL suites and cross-version runtime tests were not run locally. Spark SQL suites and macOS were skipped in CI. Disposable test code was removed, and the checkout is clean.

@andygrove andygrove added the regression A bug that did not affect the most recent Comet release label Oct 1, 2026
@parthchandra

Copy link
Copy Markdown
Contributor

The positional struct-field reconciliation here looks right, and it now makes IF more correct than CASE WHEN / COALESCE. Those two still match struct fields by name: create_case_when (native/spark-expr/src/conditional_funcs/case_when.rs:48, used by both CASE and COALESCE) computes its common type via DataFusion's get_coerce_type_for_case_expression, which folds with type_union_coercion - the same name-based matcher you just replaced with the positional if_common_type for IF.

So a case-variant struct whose field names are swapped by position, e.g. CASE WHEN q THEN named_struct('x', i, 'X', dbl) ELSE named_struct('X', int, 'x', dbl) END, would mis-coerce in Comet while Spark aligns positionally (Spark's findTypeForComplex zips fields by index and takes the left branch's name, TypeCoercionHelper.scala:238-251). The one CASE WHEN test here (if_nested_nullability.sql:113) only uses a same-name single-field struct, so it doesn't catch this.

Pre-existing and out of scope for this PR, but since this PR leaves IF ahead of CASE/COALESCE, could you file a tracking issue to extend the positional helper to the CASE/COALESCE path?

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Re-reviewed ee1b213a2f1677510d20d86a7d562ba64e169d46 against b38a60b37770df6efb4c66faac751160e92990bf with five independent scopes. No new actionable issue found; my existing approval still stands. Positional reconciliation fixes the earlier case-variant field issue, and the shared cast helper preserves CASE WHEN behavior and lazy evaluation.

Additional local validation passed on the pinned head: cargo build --locked, all 19 conditional Rust tests, and 25 Spark 4.1.3 tests across the conditional SQL files, folded-map-literal regression, IF, CASE WHEN and COALESCE suites. This includes the native JSON checks for struct names/numeric types, nested arrays/maps/structs, and mixed/all-THEN/all-ELSE batches.

The main Comet CI run passed; a separate label-triggered run has a startup failure. Local validation was focused; I did not rerun the full upstream Spark SQL matrix or other Spark profiles. CASE/COALESCE's pre-existing positional-coercion gap remains outside this change and is already tracked in #6482.

@andygrove

Copy link
Copy Markdown
Member Author

@parthchandra Filed #6532 to track extending the positional struct reconciliation to the shared CASE WHEN / COALESCE path, including the case-variant field permutation, schema/value preservation, and nested array/map coverage. I kept it separate from #6482, which tracks the simpler ELSE-branch field-name symptom and is already assigned.

@andygrove
andygrove added this pull request to the merge queue Oct 1, 2026
Merged via the queue into apache:main with commit 69c6354 Oct 2, 2026
39 checks passed
andygrove added a commit that referenced this pull request Oct 2, 2026
…6491)

* fix: cast native IF branches to a common type like CASE WHEN

Spark adds no cast to IF when its branches differ only in the nullability
of a nested field, and the native IfExpr returns a branch's array unchanged
when every row of a batch takes that branch, so the projection failed with
"column types must match schema types". Cast a branch whose Arrow type
differs from the common type, as the planner already does for CASE WHEN.

(cherry picked from commit 7645048)

* fix: adapt IF branch coercion to release conditional module

* fix: reconcile native IF struct fields positionally
grorge123 added a commit to grorge123/datafusion-comet that referenced this pull request Oct 2, 2026
…null guards

Kernel gates and null guards:
- ArraysZip guards NULL inputs with one WHEN per argument
  (`CASE WHEN a IS NULL THEN NULL WHEN b IS NULL THEN NULL ELSE arrays_zip(...) END`)
  instead of a single AND of IsNotNull, so a later argument is not evaluated on
  rows Spark's generated code skips; an ANSI division by zero there used to fail
  the query. ArraysZip with no arguments falls back to Spark, and so does
  ArraysZip over several arguments where Spark evaluates it through its
  interpreted eval, which evaluates every argument: codegen factory mode
  NO_CODEGEN, beneath a non-leaf CodegenFallback expression (array_compact's
  filter, for one), the input of an imperative aggregate (not its FILTER), or a
  generator Spark leaves out of a whole-stage stage, all of which QueryPlanSerde now tracks. In those contexts
  the codegen dispatcher compiles its kernel over the expression's eval, and
  array_repeat with a nullable count is dispatched (its eval skips the element
  when the count is NULL; the native kernel evaluates both).
- MapFromArrays checks its null guard before the LAST_WIN opt-in branch, and no
  longer refuses a literal array beside a per-row one (the native map expands the
  scalar).
- Coalesce guards only the arguments it serializes twice, and a one-argument
  coalesce (left by the optimizer only when NullPropagation and
  SimplifyConditionals are excluded) serializes as its argument instead of a CASE
  with no WHEN clause; ArrayAppend guards the item only when the array can be NULL.
- Drop the NullType CASE gate on If/CaseWhen/Coalesce (native CASE merges NullType
  branches on the locked DataFusion) and the NullType gate on collect_set (only the
  element order differed). collect_list keeps a NullType gate, documented as a guard
  against a pre-existing ungrouped final-merge failure.
- CreateArray sends a NullType argument through the dispatcher unless it is a
  literal: a foldable argument that constant folding left alone (e.g.
  `element_at(array(NULL), 1)`) is still evaluated per row, and make_array then
  returns one row for the whole batch.
- Size, ArrayAppend, ArraysZip, CollectList and CollectSet publish every reason
  they return through getUnsupportedReasons.
- CASE WHEN and coalesce cast each branch of a struct-bearing type to the
  expression's own type, made deeply nullable: a native CASE names a merged
  struct after its ELSE branch where Spark uses the first, so `array(...)`,
  array_append and the set ops met differently named structs. (Native IF
  reconciles its branches itself since apache#6458.)

Codegen dispatcher:
- The dispatcher declares its output deep-nullable (Arrow field and serialized
  return type; map keys stay non-null), as native constructors and Parquet
  columns do, so kernels that compare their inputs' types (set ops, list
  compare, CASE) see one type whatever produced each side.
- A kernel whose generated code does not compile falls back to the expression's
  interpreted eval, as Spark's whole-stage codegen does, instead of failing the
  query (e.g. element_at on an in-bounds literal index with ANSI off).
- Each occurrence of a dispatched tree that holds a Catalyst Nondeterministic
  node gets its own kernel (DispatchOccurrence): identical occurrences
  serialized to the same bytes and shared one kernel's state, so two identical
  coalesce(IF(id >= 0, monotonically_increasing_id(), NULL), -1) columns, which
  the null guard sends to the dispatcher, counted on from each other.
- The dispatcher refuses a non-deterministic user function (ScalaUDF, Invoke,
  StaticInvoke, ApplyFunctionExpression) and reflect, whose static target can
  hold JVM-wide state: Spark keeps its state in one object every call shares
  and advances it row by row, which per-expression batch kernels cannot
  reproduce. It also refuses a subquery inside the tree: the plan is serialized
  before the operator starts its subqueries, so the shipped copy failed with
  "has not finished".

Planner and native:
- arrays_overlap expands a scalar side to the batch length; it compared only the
  first row, or read past the one-row list, when a literal array met a column.
- A nested list literal whose inner lists are all empty (a folded
  `array(cast(array() AS array<bigint>))`) is built at its declared depth; it
  came out one list level deeper, and consumers that compare their inputs'
  types (set ops, arrays_overlap, CASE, =) failed. An empty child is now built
  the way a populated one is, so a deeper literal mixing the two
  (`[[[]], [[[1]]]]`) no longer fails to concatenate them.
- Set ops cast both sides to the set op's own element type, made deeply
  nullable, instead of each side's own, and cast a side even when its Spark
  type already is that type: Spark accepts struct sides whose field names
  differ only in case, a native CASE names a merged struct's fields after its
  ELSE branch, and the kernel compares field names.
- A cast that only relabels nested field names and nullability is Compatible,
  so those set-op casts stay native (array<array<date>> used to dispatch).
- Set ops lift a reason recorded under the cast they add onto themselves, and a
  Cast over a literal lifts the folded literal's reason, so both reach the
  fallback roll-up.
- Native round-robin shuffle falls back when a column it hashes contains
  NullType (the native hasher has no NullType arm); positional placement and
  columns past maxHashColumns are not hashed.
- trunc / date_trunc with a non-literal format run through the dispatcher where
  Spark evaluates them through eval, which returns NULL for an invalid format
  without evaluating the value; the check precedes every opt-in branch.
- The non-AQE DPP rewrite keeps Spark's broadcast when the broadcast gate refuses
  the build side, so the subquery still reuses the join's exchange.
- The shuffle writer rebuilds NullType builders after a batch is written, and only
  when another batch follows.

Tests and docs:
- ANSI fixtures for arrays_zip and for map_from_arrays over values from a
  dispatched lambda; a test for a one-argument coalesce kept by the optimizer;
  SQL witnesses for size, array over a non-literal foldable NullType argument,
  array_append (Spark 3.x), array_union (sides of one Spark type whose native
  nullability differs), array_intersect, coalesce, if, case when,
  collect_list, collect_set, map_from_arrays under LAST_WIN and IF/CASE branches
  of one Spark type with different native nullability, arrays_overlap of a literal
  beside a column, and a dispatched map whose generated code does not compile; the gate test covers
  Size under legacy sizeOfNull=false and the LAST_WIN order.
- A native unit test for nested list literals with empty children; a Scala test for folded nested array literals whose inner
  arrays are empty; SQL witnesses for IF / CASE / coalesce branches and set-op
  sides whose struct field names differ only in case; a CometNativeCastSuite
  test for relabel-only casts.
- Fixtures for each interpreted context (NO_CODEGEN, CodegenFallback ancestor,
  aggregate input and FILTER, generator, dispatched tree), trunc's format, and the
  round-robin NullType fallback; a CometCodegenSuite test for the compile retry.
- The composition sweep credits a case as native only when its whole plan ran in
  Comet, and keeps SimplifyExtractValueOps from removing its field extractions.
- The multi-batch shuffle test writes enough rows per destination to reach a second
  writer batch; the non-AQE DPP test checks that the build side is Comet native; UtilsSuite bounds the coalesce allocator; the composition sweep sorts
  collect_set before comparing.
- CometCodegenSuite tests for per-occurrence kernel state, non-deterministic user
  functions and subqueries inside a dispatched tree; a date_trunc interpreted
  fixture; the hash NullType witness split per function; the positional
  round-robin witness uses monotonically_increasing_id instead of a
  non-deterministic UDF.
- expressions.md, the compatibility guide, the Scala UDF guide, the shuffle guides, GenerateDocs, the expression audits and
  the benchmark describe the current routing and scan batching; expressions.md
  notes that the casts Comet adds itself follow spark.comet.expression.Cast.enabled.

Assisted-by: Claude Code (claude-opus-5-5)
grorge123 added a commit to grorge123/datafusion-comet that referenced this pull request Oct 3, 2026
…null guards

Kernel gates and null guards:
- ArraysZip guards NULL inputs with one WHEN per argument
  (`CASE WHEN a IS NULL THEN NULL WHEN b IS NULL THEN NULL ELSE arrays_zip(...) END`)
  instead of a single AND of IsNotNull, so a later argument is not evaluated on
  rows Spark's generated code skips; an ANSI division by zero there used to fail
  the query. ArraysZip with no arguments falls back to Spark, and so does
  ArraysZip over several arguments where Spark evaluates it through its
  interpreted eval, which evaluates every argument: codegen factory mode
  NO_CODEGEN, beneath a non-leaf CodegenFallback expression (array_compact's
  filter, for one), the input of an imperative aggregate (not its FILTER), or a
  generator Spark leaves out of a whole-stage stage, all of which QueryPlanSerde now tracks. In those contexts
  the codegen dispatcher compiles its kernel over the expression's eval, and
  array_repeat with a nullable count is dispatched (its eval skips the element
  when the count is NULL; the native kernel evaluates both).
- MapFromArrays checks its null guard before the LAST_WIN opt-in branch, and no
  longer refuses a literal array beside a per-row one (the native map expands the
  scalar).
- Coalesce guards only the arguments it serializes twice, and a one-argument
  coalesce (left by the optimizer only when NullPropagation and
  SimplifyConditionals are excluded) serializes as its argument instead of a CASE
  with no WHEN clause; ArrayAppend guards the item only when the array can be NULL.
- Drop the NullType CASE gate on If/CaseWhen/Coalesce (native CASE merges NullType
  branches on the locked DataFusion) and the NullType gate on collect_set (only the
  element order differed). collect_list keeps a NullType gate, documented as a guard
  against a pre-existing ungrouped final-merge failure.
- CreateArray sends a NullType argument through the dispatcher unless it is a
  literal: a foldable argument that constant folding left alone (e.g.
  `element_at(array(NULL), 1)`) is still evaluated per row, and make_array then
  returns one row for the whole batch.
- Size, ArrayAppend, ArraysZip, CollectList and CollectSet publish every reason
  they return through getUnsupportedReasons.
- CASE WHEN and coalesce cast each branch of a struct-bearing type to the
  expression's own type, made deeply nullable: a native CASE names a merged
  struct after its ELSE branch where Spark uses the first, so `array(...)`,
  array_append and the set ops met differently named structs. (Native IF
  reconciles its branches itself since apache#6458.)

Codegen dispatcher:
- The dispatcher declares its output deep-nullable (Arrow field and serialized
  return type; map keys stay non-null), as native constructors and Parquet
  columns do, so kernels that compare their inputs' types (set ops, list
  compare, CASE) see one type whatever produced each side.
- A kernel whose generated code does not compile falls back to the expression's
  interpreted eval, as Spark's whole-stage codegen does, instead of failing the
  query (e.g. element_at on an in-bounds literal index with ANSI off).
- Each occurrence of a dispatched tree that holds a Catalyst Nondeterministic
  node gets its own kernel (DispatchOccurrence): identical occurrences
  serialized to the same bytes and shared one kernel's state, so two identical
  coalesce(IF(id >= 0, monotonically_increasing_id(), NULL), -1) columns, which
  the null guard sends to the dispatcher, counted on from each other.
- The dispatcher refuses a non-deterministic user function (ScalaUDF, Invoke,
  StaticInvoke, ApplyFunctionExpression) and reflect, whose static target can
  hold JVM-wide state: Spark keeps its state in one object every call shares
  and advances it row by row, which per-expression batch kernels cannot
  reproduce. It also refuses a subquery inside the tree: the plan is serialized
  before the operator starts its subqueries, so the shipped copy failed with
  "has not finished".

Planner and native:
- arrays_overlap expands a scalar side to the batch length; it compared only the
  first row, or read past the one-row list, when a literal array met a column.
- A nested list literal whose inner lists are all empty (a folded
  `array(cast(array() AS array<bigint>))`) is built at its declared depth; it
  came out one list level deeper, and consumers that compare their inputs'
  types (set ops, arrays_overlap, CASE, =) failed. An empty child is now built
  the way a populated one is, so a deeper literal mixing the two
  (`[[[]], [[[1]]]]`) no longer fails to concatenate them.
- Set ops cast both sides to the set op's own element type, made deeply
  nullable, instead of each side's own, and cast a side even when its Spark
  type already is that type: Spark accepts struct sides whose field names
  differ only in case, a native CASE names a merged struct's fields after its
  ELSE branch, and the kernel compares field names.
- A cast that only relabels nested field names and nullability is Compatible,
  so those set-op casts stay native (array<array<date>> used to dispatch).
- Set ops lift a reason recorded under the cast they add onto themselves, and a
  Cast over a literal lifts the folded literal's reason, so both reach the
  fallback roll-up.
- Native round-robin shuffle falls back when a column it hashes contains
  NullType (the native hasher has no NullType arm); positional placement and
  columns past maxHashColumns are not hashed.
- trunc / date_trunc with a non-literal format run through the dispatcher where
  Spark evaluates them through eval, which returns NULL for an invalid format
  without evaluating the value; the check precedes every opt-in branch.
- The non-AQE DPP rewrite keeps Spark's broadcast when the broadcast gate refuses
  the build side, so the subquery still reuses the join's exchange.
- The shuffle writer rebuilds NullType builders after a batch is written, and only
  when another batch follows.

Tests and docs:
- ANSI fixtures for arrays_zip and for map_from_arrays over values from a
  dispatched lambda; a test for a one-argument coalesce kept by the optimizer;
  SQL witnesses for size, array over a non-literal foldable NullType argument,
  array_append (Spark 3.x), array_union (sides of one Spark type whose native
  nullability differs), array_intersect, coalesce, if, case when,
  collect_list, collect_set, map_from_arrays under LAST_WIN and IF/CASE branches
  of one Spark type with different native nullability, arrays_overlap of a literal
  beside a column, and a dispatched map whose generated code does not compile; the gate test covers
  Size under legacy sizeOfNull=false and the LAST_WIN order.
- A native unit test for nested list literals with empty children; a Scala test for folded nested array literals whose inner
  arrays are empty; SQL witnesses for IF / CASE / coalesce branches and set-op
  sides whose struct field names differ only in case; a CometNativeCastSuite
  test for relabel-only casts.
- Fixtures for each interpreted context (NO_CODEGEN, CodegenFallback ancestor,
  aggregate input and FILTER, generator, dispatched tree), trunc's format, and the
  round-robin NullType fallback; a CometCodegenSuite test for the compile retry.
- The composition sweep credits a case as native only when its whole plan ran in
  Comet, and keeps SimplifyExtractValueOps from removing its field extractions.
- The multi-batch shuffle test writes enough rows per destination to reach a second
  writer batch; the non-AQE DPP test checks that the build side is Comet native; UtilsSuite bounds the coalesce allocator; the composition sweep sorts
  collect_set before comparing.
- CometCodegenSuite tests for per-occurrence kernel state, non-deterministic user
  functions and subqueries inside a dispatched tree; a date_trunc interpreted
  fixture; the hash NullType witness split per function; the positional
  round-robin witness uses monotonically_increasing_id instead of a
  non-deterministic UDF.
- expressions.md, the compatibility guide, the Scala UDF guide, the shuffle guides, GenerateDocs, the expression audits and
  the benchmark describe the current routing and scan batching; expressions.md
  notes that the casts Comet adds itself follow spark.comet.expression.Cast.enabled.

Assisted-by: Claude Code (claude-opus-5-5)
grorge123 added a commit to grorge123/datafusion-comet that referenced this pull request Oct 4, 2026
…null guards

Kernel gates and null guards:
- ArraysZip guards NULL inputs with one WHEN per argument
  (`CASE WHEN a IS NULL THEN NULL WHEN b IS NULL THEN NULL ELSE arrays_zip(...) END`)
  instead of a single AND of IsNotNull, so a later argument is not evaluated on
  rows Spark's generated code skips; an ANSI division by zero there used to fail
  the query. ArraysZip with no arguments falls back to Spark, and so does
  ArraysZip over several arguments where Spark evaluates it through its
  interpreted eval, which evaluates every argument: codegen factory mode
  NO_CODEGEN, beneath a non-leaf CodegenFallback expression (array_compact's
  filter, for one), the input of an imperative aggregate (not its FILTER), or a
  generator Spark leaves out of a whole-stage stage, all of which QueryPlanSerde now tracks. In those contexts
  the codegen dispatcher compiles its kernel over the expression's eval, and
  array_repeat with a nullable count is dispatched (its eval skips the element
  when the count is NULL; the native kernel evaluates both).
- MapFromArrays checks its null guard before the LAST_WIN opt-in branch, and no
  longer refuses a literal array beside a per-row one (the native map expands the
  scalar).
- Coalesce guards only the arguments it serializes twice, and a one-argument
  coalesce (left by the optimizer only when NullPropagation and
  SimplifyConditionals are excluded) serializes as its argument instead of a CASE
  with no WHEN clause; ArrayAppend guards the item only when the array can be NULL.
- Drop the NullType CASE gate on If/CaseWhen/Coalesce (native CASE merges NullType
  branches on the locked DataFusion) and the NullType gate on collect_set (only the
  element order differed). collect_list keeps a NullType gate, documented as a guard
  against a pre-existing ungrouped final-merge failure.
- CreateArray sends a NullType argument through the dispatcher unless it is a
  literal: a foldable argument that constant folding left alone (e.g.
  `element_at(array(NULL), 1)`) is still evaluated per row, and make_array then
  returns one row for the whole batch.
- Size, ArrayAppend, ArraysZip, CollectList and CollectSet publish every reason
  they return through getUnsupportedReasons.
- CASE WHEN and coalesce cast each branch of a struct-bearing type to the
  expression's own type, made deeply nullable: a native CASE names a merged
  struct after its ELSE branch where Spark uses the first, so `array(...)`,
  array_append and the set ops met differently named structs. (Native IF
  reconciles its branches itself since apache#6458.)

Codegen dispatcher:
- The dispatcher declares its output deep-nullable (Arrow field and serialized
  return type; map keys stay non-null), as native constructors and Parquet
  columns do, so kernels that compare their inputs' types (set ops, list
  compare, CASE) see one type whatever produced each side.
- A kernel whose generated code does not compile falls back to the expression's
  interpreted eval, as Spark's whole-stage codegen does, instead of failing the
  query (e.g. element_at on an in-bounds literal index with ANSI off).
- Each occurrence of a dispatched tree that holds a Catalyst Nondeterministic
  node gets its own kernel (DispatchOccurrence): identical occurrences
  serialized to the same bytes and shared one kernel's state, so two identical
  coalesce(IF(id >= 0, monotonically_increasing_id(), NULL), -1) columns, which
  the null guard sends to the dispatcher, counted on from each other.
- The dispatcher refuses a non-deterministic user function (ScalaUDF, Invoke,
  StaticInvoke, ApplyFunctionExpression) and reflect, whose static target can
  hold JVM-wide state: Spark keeps its state in one object every call shares
  and advances it row by row, which per-expression batch kernels cannot
  reproduce. It also refuses a subquery inside the tree: the plan is serialized
  before the operator starts its subqueries, so the shipped copy failed with
  "has not finished".

Planner and native:
- arrays_overlap expands a scalar side to the batch length; it compared only the
  first row, or read past the one-row list, when a literal array met a column.
- Tests for a nested list literal whose inner lists are all empty (a folded
  `array(cast(array() AS array<bigint>))`), and for a deeper literal mixing
  empty and populated children (`[[[]], [[[1]]]]`). Main now builds both at
  their declared depth (apache#4715); consumers that compare their inputs' types
  (set ops, arrays_overlap, CASE, =) failed on the old one-level-deeper build.
- Set ops cast both sides to the set op's own element type, made deeply
  nullable, instead of each side's own, and cast a side even when its Spark
  type already is that type: Spark accepts struct sides whose field names
  differ only in case, a native CASE names a merged struct's fields after its
  ELSE branch, and the kernel compares field names.
- A cast that only relabels nested field names and nullability is Compatible,
  so those set-op casts stay native (array<array<date>> used to dispatch).
- Set ops lift a reason recorded under the cast they add onto themselves, and a
  Cast over a literal lifts the folded literal's reason, so both reach the
  fallback roll-up.
- Native round-robin shuffle falls back when a column it hashes contains
  NullType (the native hasher has no NullType arm); positional placement and
  columns past maxHashColumns are not hashed.
- trunc / date_trunc with a non-literal format run through the dispatcher where
  Spark evaluates them through eval, which returns NULL for an invalid format
  without evaluating the value; the check precedes every opt-in branch.
- The non-AQE DPP rewrite keeps Spark's broadcast when the broadcast gate refuses
  the build side, so the subquery still reuses the join's exchange.
- The shuffle writer rebuilds NullType builders after a batch is written, and only
  when another batch follows.

Tests and docs:
- ANSI fixtures for arrays_zip and for map_from_arrays over values from a
  dispatched lambda; a test for a one-argument coalesce kept by the optimizer;
  SQL witnesses for size, array over a non-literal foldable NullType argument,
  array_append (Spark 3.x), array_union (sides of one Spark type whose native
  nullability differs), array_intersect, coalesce, if, case when,
  collect_list, collect_set, map_from_arrays under LAST_WIN and IF/CASE branches
  of one Spark type with different native nullability, arrays_overlap of a literal
  beside a column, and a dispatched map whose generated code does not compile; the gate test covers
  Size under legacy sizeOfNull=false and the LAST_WIN order.
- A native unit test for nested list literals with empty children; a Scala test for folded nested array literals whose inner
  arrays are empty; SQL witnesses for IF / CASE / coalesce branches and set-op
  sides whose struct field names differ only in case; a CometNativeCastSuite
  test for relabel-only casts.
- Fixtures for each interpreted context (NO_CODEGEN, CodegenFallback ancestor,
  aggregate input and FILTER, generator, dispatched tree), trunc's format, and the
  round-robin NullType fallback; a CometCodegenSuite test for the compile retry.
- The composition sweep credits a case as native only when its whole plan ran in
  Comet, and keeps SimplifyExtractValueOps from removing its field extractions.
- The multi-batch shuffle test writes enough rows per destination to reach a second
  writer batch; the non-AQE DPP test checks that the build side is Comet native; UtilsSuite bounds the coalesce allocator; the composition sweep sorts
  collect_set before comparing.
- CometCodegenSuite tests for per-occurrence kernel state, non-deterministic user
  functions and subqueries inside a dispatched tree; a date_trunc interpreted
  fixture; the hash NullType witness split per function; the positional
  round-robin witness uses monotonically_increasing_id instead of a
  non-deterministic UDF.
- expressions.md, the compatibility guide, the Scala UDF guide, the shuffle guides, GenerateDocs, the expression audits and
  the benchmark describe the current routing and scan batching; expressions.md
  notes that the casts Comet adds itself follow spark.comet.expression.Cast.enabled.

Assisted-by: Claude Code (claude-opus-5-5)
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 5, 2026
The planner promised DataFusion the kernel's own return type, with only list
and map child names canonicalized. A kernel may report a nested field as
non-nullable where the registration declares it nullable: the declared-vs-
reported check disregards nested nullability, so this is accepted. An operator
that takes its type from the UDF would then have to narrow a nullable field
from another expression to match, and that cast fails with `Cannot cast
nullable struct field 'a' to non-nullable field`.

`canonicalize_child_names` becomes `promised_return_type`, which also makes
every nested field nullable except a map's entries and keys, which Arrow
requires to be non-null. The adapter accepts a result whose promised form
matches and conforms it with a cast that only renames fields and widens their
nullability, reusing the buffers.

`make_struct_c` in the test library reports a non-nullable struct field. The
planner and adapter tests pin the widened type. The suite runs the reported
`if` query in both branch orders. Main's IF branch reconciliation (apache#6458) now
covers that query too, so the planner change keeps other consumers consistent.
LinSimon-901101 pushed a commit to LinSimon-901101/datafusion-comet that referenced this pull request Oct 6, 2026
* feat: custom Rust UDFs via arrow-ffi (alternative to apache#4283) [skip ci]

Adds custom Rust scalar UDF support using arrow-only FFI surfaces, as
suggested by paleolimbot and timsaucer in feedback on apache#4283. This is an
alternative implementation for comparison; see apache#4283 for the bespoke-ABI
version.

Two ABI flavors are provided side-by-side so reviewers can compare:

  1. C ABI (sedona-style): pure C-callable struct of function pointers,
     parameterized only by Arrow C Data Interface (FFI_ArrowSchema /
     FFI_ArrowArray). Decoupled from datafusion versions; future-portable
     to C/C++. Modeled on apache/sedona-db's SedonaCScalarKernel header.

  2. datafusion-ffi (FFI_ScalarUDF): wraps user's ScalarUDFImpl as
     FFI_ScalarUDF. Inherits full ScalarUDFImpl surface (variadic
     signatures, type coercion, metadata-aware return types) for free,
     at the cost of a major-version pin against datafusion-ffi.

A single library may export either or both; loader walks both discovery
functions. `comet-test-udfs` exposes `add_one_c` (C ABI) and `add_one_df`
(datafusion-ffi) and the e2e suite drives both through Spark.

Scope is intentionally minimal for comparison: scalar-only, happy-path
e2e tests only (no panic/error/signature-mismatch coverage). The Scala
JVM API mirrors apache#4283 exactly so the comparison is apples-to-apples.

Pieces:
- native/comet-udf-sdk: SDK with both ABI flavors + export macros
- native/comet-test-udfs: cdylib exposing one UDF per ABI
- native/core/src/execution/rust_udf: loader, cache, ImportedCScalarUdf
- native/core/src/comet_rust_udf_bridge.rs: JNI for validateLibrary/listUdfs
- native/proto: RustUdfCall message
- spark/.../udf: CometRustUDF.register, registry, exception classes,
  JNI bridge stub
- spark/.../serde/CometScalaUDF: dispatch ScalaUDF to RustUdfCall
  when udfName is in the registry

* feat: drop datafusion-ffi flavor, contain panics at the FFI boundary, add user guide

Keep a single UDF ABI built only on the Arrow C Data Interface. Carrying two
public ABIs would mean two permanent compatibility promises; the datafusion-ffi
flavor pins every user cdylib to Comet's DataFusion major version, which the
53 -> 54 upgrade demonstrated by removing as_any from ScalarUDFImpl. The
rationale is recorded in comet-udf-sdk's crate docs.

Contain panics at every extern "C" entry point. A panic escaping one aborts
the process, taking down the executor JVM and every task on it. User UDF code
is arbitrary and panicking is idiomatic Rust, so a panic is now converted to a
query error via the existing get_last_error channel. Previously only execute
was guarded; init, new_impl, the discovery entry point and the release
callbacks were not.

Add a user guide page describing the feature as experimental, including the
executor-side library placement requirement and the trust implications of
loading native code.

* refactor: remove the write-only spark.comet.rustUdfs conf

propagateConf wrote the registered UDF set into spark.comet.rustUdfs on every
register call, and nothing ever read it. Executors resolve the library from the
library_path carried in the RustUdfCall proto, so the conf implied a
propagation mechanism that does not exist. Document what actually happens on
the register scaladoc instead. CometRustUdfRegistry.snapshot had no other
caller and goes with it.

* feat: cover all non-nested Spark types, validate the declared return type

Add echo_c (identity over any type) and stringify_c ((any) -> Utf8) to the test
cdylib, and drive both over every non-nested Spark type with a null row: bool,
the four int widths, float, double, decimal, string, binary, date, timestamp and
timestamp_ntz. echo_c proves the array survives the round trip with its type and
nulls; stringify_c forces the UDF to decode the values rather than hand the
array straight back.

No ABI change was needed, since the FFI surface is type-agnostic, but nothing
had demonstrated that beyond Int64.

Also validate at plan time that the return type declared through
CometRustUDF.register matches what the kernel's return_field reports, naming
both types and pointing at register. Registering decimal(10,2) for a column
Spark had widened to decimal(11,2) previously surfaced as a bare DataFusion
type assertion partway through execution.

Complex types (array, struct, map) remain future work and are documented as
unsupported.

* feat: support complex types, compare return types ignoring nested nullability

Arrays, maps and structs work over the ABI, including nesting (array of struct,
struct of array), so extend the type matrix to cover them alongside the
non-nested types. No ABI change was needed; the FFI surface carries child arrays
already.

The new plan-time return type check rejected them, and was right to look: Spark
carries containsNull and field nullability inside the declared type, so
array<int> with containsNull=false converts to a List with a non-nullable child,
while the array actually delivered to the UDF has that child normalized to
nullable. Compare with nested nullability erased, and promise DataFusion the
kernel's own type rather than the declared one so its exact-match assertion
compares like with like. The comparison stays strict about everything that
changes how bytes are read: decimal precision and scale, timestamp unit, struct
field names and order.

Also document the return type model, which came up in review. A kernel has no
fixed return type: return_field derives one from the argument types on every
call, and echo_c serves all 19 types in the matrix. What is fixed is the type
declared to register, because Spark needs a concrete DataType to plan against,
and that is per-registration rather than per-kernel.

* ci: run CometRustUdfSuite in the PR builds

The suite cancelled itself whenever -Dcomet.test.udfs.lib was unset, and nothing
set it, so it had never run in CI and the ABI had never been exercised against a
Linux .so.

Have the suite locate the cdylib under native/target instead of requiring the
flag, keeping the property as an override. When the library is missing the suite
still skips locally, but fails in CI, where its absence means the native build or
the artifact upload changed rather than that someone forgot a flag. An undefined
Maven property arrives in the forked JVM as the literal string "null", which is
why the override needs filtering rather than a plain null check.

Upload libcomet_test_udfs alongside libcomet from both native build jobs: the
same cargo invocation already produces it, since the crate is a workspace
default member, and the test jobs download the artifact rather than building.

Register the suite in the expressions bucket of both workflows, as
dev/ci/check-suites.py requires.

* style: drop redundant string interpolators in the Rust UDF messages

scalafix RedundantSyntax flags s-prefixed literals with nothing to
interpolate. These were latent: the previous head commit on this branch carried
[skip ci], so the lint jobs had never run against these files.

* fix: drop a loaded UDF library after its kernels, not before

LoadedLibrary declared `library: Arc<Library>` ahead of `udfs`. Struct
fields drop in declaration order, so dropping a LoadedLibrary ran dlclose
first and then invoked each kernel's `release` function pointer, which
lives in the text of the library just unloaded. Nothing else holds a clone
of that Arc.

Unreachable in production, since the process-wide cache never drops a
LoadedLibrary, but the loader tests build one via `load()` and drop it.

Reorder the fields and say why in a comment, so a later edit does not
quietly reintroduce it.

* fix(ci): build the test UDF cdylib before the Rust test run

The rust-test job runs `cargo nextest run` and nothing else. comet-test-udfs
is `crate-type = ["cdylib"]` with no test targets, so a test build compiles
the crate without ever emitting libcomet_test_udfs, and the three rust_udf
tests that dlopen it failed on a clean checkout:

    execution::rust_udf::cache::tests::same_path_returns_same_arc
    execution::rust_udf::loader::tests::all_exported_kernels_are_discovered
    execution::rust_udf::loader::tests::load_test_udfs_succeeds

It passed locally only because a previous `cargo build` had left the
artifact in target/debug.

Build the crate explicitly before the test run, and replace the bare
`expect("load")` with a message naming that command, so the next person to
hit this reads the fix instead of a dlopen error.

* fix: reject deterministic = false when registering a Rust UDF

`CometRustUDF.register` accepted a `deterministic` flag, installed a
nondeterministic Spark stub for it, and carried it in the RustUdfCall proto,
but nothing on the native side ever read it: ImportedCScalarUdf hardcodes
Volatility::Immutable. A UDF the caller declared nondeterministic was
therefore planned as pure and could be constant-folded, evaluated once and
reused, or eliminated as a common subexpression.

Honoring the flag needs the cached per-library signature to become per
registration, which is apache#5249. Until then, fail the registration rather than
let the parameter lie about what Comet does with it.

Raised in review: does the proto field's value mean RustUdfs are always
immutable?

* docs: clarify the Rust UDF ABI's stability, purity, and path resolution

Review follow-ups on the SDK and the user guide, plus two assertion tweaks:

- State that the C ABI structs are specific to one Comet version and are
  internal under the versioning policy, while the host's own arrow and
  datafusion versions need not match the cdylib's.
- Say on CometCScalarUdf that only immutable functions are supported, and
  add the same to the user guide's limitations. The guide previously told
  readers to pass `deterministic = false` for impure functions, which is now
  refused.
- Describe what libraryPath actually does: an absolute path is the sensible
  choice, but a bare name resolves through the platform loader search path,
  and neither is a security boundary.
- Note that CometRustUDF is deliberately not @public, so it carries no
  compatibility guarantee.
- Report rather than swallow a panic caught in a release callback: there is
  no error channel there, so it goes to stderr with a context string.
- Check `release.is_some()` alongside private_data in the FFI debug asserts,
  which catches a call made after release rather than only an uninitialized
  one.

* review: address mbutrovich's first pass on the Rust UDF path

Drop the dead code and vacuous round-trip the review found, correct the
comments that were wrong, and pin the error classification with tests.

- `validateLibrary` returns void. It only ever reported back the name it
  was given, since the native side already filters `lib.udfs` by that
  name and errors when it is absent, so the JSON construct/parse cycle
  and the `require(described.name == name)` that followed it could not
  fail. `listUdfs` had no caller and comes out with it.
- `classifyNativeError` matched a bare "ABI", which every loader message
  can contain because they all interpolate the library path: a library
  under a directory named `ABI` had its "failed to open" reported as an
  ABI mismatch. Match the message wording instead, and add a test per
  failure mode so a rewording on the native side fails a test rather
  than silently changing the exception a caller sees.
- Publish the registry entry before installing the catalog stub. The
  stub is what makes the name resolvable to the analyzer, so the old
  order left a window where a query could plan the name with no registry
  entry and hit the stub's "not evaluated" exception.
- Correct the `Mutex` comment in `ImportedCScalarUdf`. It claimed
  DataFusion serializes invocations of a `ScalarUDFImpl` per batch; it
  does not, and the planner hands every task the same instance out of
  the process-wide cache, so concurrent calls are the normal case. The
  lock is load-bearing rather than defensive, and it serializes all
  batches of a UDF in a process. Relevant to apache#5252, which would remove
  it.
- Document `scalar_args` as reserved and always NULL in ABI v1. Nothing
  on the host can populate it: literals are expanded to full-length
  arrays before a kernel sees them.
- Justify the `CometRustUdfRegistry` singleton's lifetime and bounds
  where the contributor guide asks for it, and record that session
  scoping is unresolved.
- Add the tests the review asked for: a concurrent-registration test, and
  a name-collision test which is `ignore`d because it fails today. A
  Rust UDF does answer calls to an ordinary Scala UDF registered under
  the same name. The obvious fix does not work: `functions.udf` wraps the
  closure it is given, so the object in `ScalaUDF.function` is Spark's
  and not one Comet can recognize, and the `udf(AnyRef, DataType)`
  overload that would have preserved identity is gone in Spark 4.

* docs: point the Rust UDF follow-up comments at their tracking issues

Replaces the review-thread links added in the previous commit, and adds
pointers where the code has a known gap but carried no reference:

- apache#5294 registry session scoping
- apache#5295 a Rust UDF answers an ordinary Scala UDF of the same name
- apache#5296 the adapter rebuilds the impl and re-resolves the return type
  per batch
- apache#5297 the library cache holds its write lock across dlopen

* refactor: name the UDF mechanism for what it is, not for Rust

The ABI this feature loads a UDF through is a set of `#[repr(C)]` structs
of function pointers carrying Arrow C Data Interface arrays. No DataFusion
type and no Rust type appears in it, the entry points are unmangled C
symbols, and every allocation is freed through a callback the library
supplies, so nothing about the host requires the library to have been
written in Rust. The Rust SDK is one client of that ABI, not the ABI
itself.

The names said otherwise, and one of them is a wire format that would
have been awkward to change later:

- proto `RustUdfCall` -> `NativeScalarUdf`, field `rust_udf_call` ->
  `native_scalar_udf`, which also parallels the existing `JvmScalarUdf`
  for the JVM codegen dispatch path. Field number 75 is unchanged.
- `CometRustUDF` -> `CometNativeUDF`, and the registry, bridge, metadata,
  exceptions and suite alongside it.
- Rust internals move to `execution::c_udf`, matching the C-ABI names
  already there (`CometCScalarKernel`, `comet_c_udf_list_v1`), and the
  bridge file to `comet_native_udf_bridge.rs`, matching the JVM class it
  serves.
- User-facing messages and comments say "native UDF" where they mean the
  mechanism, and keep "Rust" where they mean the SDK.

Unchanged on purpose: `comet-udf-sdk` and `comet-test-udfs` were already
neutral, and `rust_udfs.md` documents the Rust SDK specifically, so a
`cpp_udfs.md` would sit beside it rather than replace it.

Also narrows the claim that went with the old naming. The user guide said
the ABI "is implementable from C or C++", which reads as an offer; it is
true of the design but nothing supports it, since no header is published,
the layouts are only defined in Rust source and are not stable across
releases, and the SDK's panic guards have no automatic equivalent in
another language. The docs now say that plainly.

Drive-by: the `NativeScalarUdf` proto comment still described dispatch
through "whichever ABI flavor (C ABI / datafusion-ffi)" the library
registered under. The datafusion-ffi flavor was removed earlier in this
branch.

* perf: drop the per-process mutex serializing native UDF batches

The adapter held one `Mutex` across the whole `invoke_with_args` body,
including the user's compute. Every Spark task in an executor shares one
`ImportedCScalarUdf` through the process-wide library cache, so a UDF over
a wide scan ran at roughly one core no matter how many tasks were in
flight.

The lock is not needed. Both kernel-level callbacks the adapter reaches,
`function_name` and `new_impl`, take `*const CometCScalarKernel`, and the
SDK's exporter only reads through that pointer: everything mutable is the
per-call `CometCScalarKernelImpl`, which never leaves the stack frame that
built it. State the concurrency requirement plainly in the ABI docs, on
the `CometCScalarUdf` trait and in the user guide, and pin it with a test
driving eight threads through one shared adapter.

* fix: stop rejecting a native UDF over positional list and map field names

Two type-mapping traps, both reachable by following the user guide.

The guide listed `TimestampType` and `TimestampNTZType` as the same Arrow
type. They are not: Comet maps the former to
`Timestamp(Microsecond, Some("UTC"))` and the latter to
`Timestamp(Microsecond, None)`, and the timezone is compared exactly, so
a UDF that built a `TimestampMicrosecondArray` without calling
`with_timezone` was rejected at planning against a `TimestampType`
registration with no hint as to which axis disagreed. The table now shows
the tag, the error message names the axes that must match exactly, and
`make_ts_utc_c` / `make_ts_naive_c` pin both directions.

The map case was a real rejection over a spelling. Comet emits a map's
children as `entries` / `key` / `value` while arrow-rs's
`MapBuilder::new(None, ..)` names them `entries` / `keys` / `values`, and
the comparison was strict about names at every level. Those names are
positional in the Arrow columnar format and Comet's own `CometMapVector`
reads them by index, so they are now normalized before comparing, as is a
list's element name. Struct field names stay strict: they are part of the
Spark type and are how a caller addresses the result.

`echo_c` could not have caught either one, since it hands back the array
Comet gave it and so inherits Comet's own tag and names. `make_map_c`
builds its output through arrow-rs's default builder instead.

* test: cover native UDFs in filter, join, grouping key and window

Every end-to-end test called the UDF from a projection. Filters, join
conditions, grouping keys and window partitioning each route through a
different Comet operator, so a regression in the rule wiring for any of
them would have shown up only as a silent fallback to Spark.

The catalog stub is what makes these assertions load-bearing: it throws
if Spark ever evaluates the UDF itself, so each test fails rather than
passing on the JVM path.

* refactor: expose register and get on the CometNativeUdfRegistry object

Callers reached through `CometNativeUdfRegistry.instance` to get at the
process-wide registry, which puts the singleton in every call site for no
benefit. Forward both methods from the companion so `instance` stays an
implementation detail.

* test: load the concurrent adapter test's library through the cache

concurrent_invocations_share_one_adapter took its adapter from a
LoadedLibrary returned by `load` and let that struct drop before the
threads ran. The adapter holds no reference to its library, so the drop
dlclosed it. On Linux that unmaps the library, and the first `new_impl`
call jumps into unmapped text. The Linux rust-test job has failed on
both pushes since the mutex was removed, and the 2026-09-07 log shows
this test dying with SIGSEGV. macOS keeps an image with thread-locals
mapped after dlclose, so it passed locally.

Go through `cache::get_or_load` instead, which is how the planner gets
an adapter and which never unloads. Also correct the `library` field
doc, which said the library is never unloaded when only the cache makes
that true.

* fix: refuse a UDF result whose type disagrees with return_field

An FFI_ArrowArray carries no type, so the host imports a result as whatever
return_field declared. A UDF declaring Int64 that returned an Int32Array was read
as the wrong values, or past the end of the buffer. Check the result in the SDK,
the last point that knows both types, and correct the crate docs that still
promised cross-release compatibility and C/C++ support.

* fix: keep native UDF libraries loaded and plan canonical child names

The adapter now holds an Arc<Library> declared after its kernel, so it cannot
outlive the library its callbacks live in. The planner promises the kernel's
return type with list and map children renamed to Comet's names, and the adapter
relabels each result to match, so a map built with MapBuilder defaults combines
with Comet's own maps in if/CASE. Adds sub_c and tests for scalar arguments.

* test: cover multi-argument, literal and if-combined native UDF calls

Also removes CometNativeUdfSignatureException, which was never thrown.

* docs: correct the Rust UDF guide's dependency, panic and reload advice

* feat: refuse native UDF calls whose argument types differ from the registered ones

The catalog stub is untyped, so Spark inserts no casts and inputTypes was only
used for its length. The serde now compares the call's argument types with the
registered ones, ignoring nullability, and fails at planning time naming both.

Closes the open question from the self-review on what inputTypes means.

* fix: release the native UDF cache's read guard before taking the write lock

`get_or_load` looked up the canonical path in an `if let` scrutinee and took
the write lock inside its body to record the new path. The scrutinee's read
guard lives until the end of that block, so a first lookup through a path that
resolves to an already-loaded library, such as a symlink, waited on its own
read lock forever. Two threads racing on a library's first load through a
non-canonical path could reach the same branch.

The lookup is now a statement of its own. The new test looks a loaded library
up through a symlink on a separate thread, so a regression fails it after 30
seconds instead of hanging the run.

* fix: widen a native UDF's nested nullability before it enters the plan

The planner promised DataFusion the kernel's own return type, with only list
and map child names canonicalized. A kernel may report a nested field as
non-nullable where the registration declares it nullable: the declared-vs-
reported check disregards nested nullability, so this is accepted. An operator
that takes its type from the UDF would then have to narrow a nullable field
from another expression to match, and that cast fails with `Cannot cast
nullable struct field 'a' to non-nullable field`.

`canonicalize_child_names` becomes `promised_return_type`, which also makes
every nested field nullable except a map's entries and keys, which Arrow
requires to be non-null. The adapter accepts a result whose promised form
matches and conforms it with a cast that only renames fields and widens their
nullability, reusing the buffers.

`make_struct_c` in the test library reports a non-nullable struct field. The
planner and adapter tests pin the widened type. The suite runs the reported
`if` query in both branch orders. Main's IF branch reconciliation (apache#6458) now
covers that query too, so the planner change keeps other consumers consistent.

* fix: ignore struct field metadata when checking native UDF argument types

`checkArgumentTypes` compared `deepNullable` forms, which keep each
`StructField`'s metadata. A struct argument whose field carries a comment was
refused against a registration with the same field name and type, with an
error naming `struct<a:int>` on both sides.

`DataTypeSupport.equalsIgnoreNullability` re-derives Spark's
`DataType.equalsIgnoreNullability`, which compares struct field names but
neither nullability nor metadata. Spark's own is `private[sql]` on 3.4. The
new test passes a struct whose field has a comment, built with
`struct(col.as(name, metadata))`.

* fix: discard native UDF registration results explicitly for the strict-warnings build

`spark.udf.register` returns the registered `UserDefinedFunction` and
`ConcurrentHashMap.put` returns the previous value. Both were discarded
implicitly at the end of a `Unit` method, which `-Ywarn-value-discard` under
`-Pstrict-warnings` rejects as `discarded non-Unit value`. The fatal error
also cut scalac short, which is where the job's four errors in untouched
shuffle files came from.

* docs: pin the Rust UDF SDK to the release the docs were built for

The example dependency pinned `tag = "1.1.0"`, which predates the SDK, so it
could not resolve. It now uses `$COMET_VERSION`, which the docs build replaces
with the release a page was built for. The development docs, whose version has
no tag, say to pin `rev` instead.

* fix: recognize native UDFs by their function instead of by name

`CometScalaUDF` recognized a native UDF by looking the call's `udfName` up in
`CometNativeUdfRegistry`, a process-wide map keyed by bare name. So an
ordinary UDF registered under a native UDF's name was answered out of the
native library, returning wrong values with no warning (apache#5295). A
registration in one session also claimed the name in every other session
sharing the driver JVM (apache#5294).

`register` now installs a builder in the session's own function registry, as
`spark.udf.register` does. The builder creates a `ScalaUDF` that holds a
`CometNativeUdfFunction`, which carries the registration, and the serde
matches on that function. An ordinary UDF registered under the same name holds
its own function and stays on the JVM path. A registration is a temporary
function of the session it was made in. The process-wide registry is gone.

`spark.udf.register` could not do this, because `functions.udf` wraps the
function it is handed, and the overload that did not is gone in Spark 4. The
builder checks the argument count itself, which the fixed-arity stub used to
do implicitly. `ScalaUDF` casts its function to the `FunctionN` of its arity,
so there is one subclass per arity, up to the existing limit of four.

The name-collision reproduction is enabled. A new test checks that another
session neither sees the registration nor has its own UDF of that name
answered natively.

* refactor: register native UDFs as session-scoped Catalyst expressions

Match the registration in apache#6697. A native UDF's calls now resolve to a
`NativeUdfCall` expression rather than a `ScalaUDF` holding a marker
function, and `CometNativeUdfCall`, a serde of its own, emits the
`NativeScalarUdf`. `CometScalaUDF` and `DataTypeSupport` are back to main's
versions.

- Spark's analyzer checks the argument types (`ExpectsInputTypes`),
  disregarding nullability and field metadata, and inserts no casts. A
  mismatch fails analysis instead of Comet's planning, so Comet's own type
  check, `equalsIgnoreNullability` and `CometNativeUdfArgumentTypeException`
  are gone.
- A call with the wrong number of arguments fails analysis with Spark's own
  error.
- Nothing caps the arity any more. Four was the most `ScalaUDF`'s `FunctionN`
  casts allowed for, so `sum_c` and a test cover a five-argument call.
- The codegen dispatcher declines an ordinary UDF whose arguments hold a
  native call, so Spark gets that operator and the failure names the real
  cause.
- If Spark evaluates the call, it throws `CometUdfNotEvaluatedException`.

`ShimSessionFunctionRegistry`, `CometUdfErrors` and `CometUdfExceptions` are
byte-identical copies of apache#6697's, so whichever pull request lands second
merges them without a conflict.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:expressions Expression evaluation backport-1.1 Candidate for backporting to 1.1 release branch bug Something isn't working regression A bug that did not affect the most recent Comet release

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Native IF fails with "column types must match schema types" when its branches differ in nested nullability

4 participants