Skip to content

fix: return an empty array from CASE WHEN for an empty batch - #6763

Open
viirya wants to merge 2 commits into
apache:mainfrom
viirya:fix-case-when-empty-batch-scalar
Open

viirya wants to merge 2 commits into
apache:mainfrom
viirya:fix-case-when-empty-batch-scalar

Conversation

@viirya

@viirya viirya commented Oct 7, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6762.

Rationale for this change

IN builds its candidate list once when every candidate is a constant. Both DataFusion's InListExpr and Comet's spark_in_list decide that by evaluating each candidate on an empty batch and checking for a scalar. CaseWhenExpr, which IF, nullif and coalesce also use, returned a scalar NULL for an empty batch, because no row chose a branch. So a CASE or IF that reads a column was taken as the constant NULL, and for example id IN (IF(id = 1, NULL, id)) returned NULL for every row.

DataFusion has the same bug in InListExpr::try_new for its own CaseExpr (apache/datafusion#26082), fixed by apache/datafusion#26083, which builds the static list only from items whose leaves are all literals. That does not replace this fix: Comet plans spark_partition_id() as a literal, so CASE WHEN spark_partition_id() = 0 THEN 1L END still has only literal leaves, and still has to return its value rather than a NULL for the empty batch.

What changes are included in this PR?

  • CaseWhenExpr::evaluate returns an empty array of the result type for an empty batch. The check runs before both the eager evaluation from perf: optimize native CASE WHEN and IF (up to 12x faster) #6350 and the lazy one, because DataFusion's CaseExpr also returns a scalar NULL there, for example for a NULL THEN literal. As a result, a CASE made only of literals is no longer treated as a constant by IN. That only affects performance, and Spark folds such a CASE before Comet sees it.
  • spark_in_list treats a candidate as a constant only when every leaf is a Literal and no node is volatile, the same rule as fix: build the IN list static filter only from literal-only items datafusion#26083, so a candidate that reads a column stays dynamic whatever it returns for the empty batch.
  • adding_a_new_expression.md says that an expression whose result depends on a column must return an array, including for an empty batch.

I also looked for other Comet expressions that return a scalar for an empty batch while reading a column, and found none. The other PhysicalExpr implementations and scalar functions return a scalar only when all their inputs are scalars, or when a NULL scalar argument makes every row NULL.

IN still evaluates later candidates over the whole batch rather than per row, which predates this PR and is tracked in #6823.

How are these changes tested?

  • Rust unit tests: CASE and IF return an empty array for an empty batch on both the eager and the lazy path, and spark_in_list keeps dynamic a candidate that reads a column but returns a scalar for the empty batch. Both fail without the fix.
  • expressions/conditional/in_case_when_candidate.sql runs the queries from the issue with native range: IF, nullif, coalesce, CASE with and without ELSE, NOT IN, a CASE over spark_partition_id(), and struct and array operands, including a double leaf so that spark_in_list is used. It fails without the fix.

This pull request and its description were written by Isaac.

IN builds its candidate list once when every candidate is a constant, and
decides that by evaluating each candidate on an empty batch and checking for
a scalar. Both DataFusion's InListExpr and Comet's spark_in_list do this.

CaseWhenExpr, which IF and nullif also use, returned a scalar NULL for an
empty batch because no row chose a branch. A CASE or IF that reads a column
was therefore taken as the constant NULL, so for example
`id IN (IF(id = 1, NULL, id))` returned NULL for every row. A named_struct
over such a CASE was wrong the same way, since it returns a scalar when all
its children do.

Return an empty array of the result type for an empty batch, before either
the eager or the lazy evaluation, since DataFusion's CaseExpr also returns a
scalar NULL there when there is no ELSE. Also make spark_in_list treat a
candidate that reads a column as non-constant, whatever it returns for the
empty batch.
@github-actions github-actions Bot added bug Something isn't working area:expressions Expression evaluation labels Oct 7, 2026
@viirya
viirya requested review from andygrove and sunchao October 7, 2026 20:35

@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.

Reviewed the complete three-file diff from b56349697b786ff2ad1c1bcf6ecf45b809af5f30 to ebfd4005f91176e3d62a969cf6ba1a360d5ffb09. The PR is not a draft. There are no existing reviews, conversation comments, inline comments, or review threads. One newly exposed P2 regression was found.

Routed skills: .ai/skills/review-comet-pr/SKILL.md and .ai/skills/review-comet-expression-pr/SKILL.md.

Summary

  • Prior state and problem: Empty-batch evaluation could turn a column-dependent CASE, IF, or nullif candidate into a cached scalar NULL. Queries such as id IN (IF(id = 1, NULL, id)) consequently returned incorrect membership results.
  • Design approach: Return a typed empty array before either CASE evaluation path, and reject column-reading candidates from spark_in_list constant detection.
  • Correctness: The intended membership correction is covered by passing tests. However, making these candidates dynamic exposes whole-batch IN evaluation: an ANSI cast can now execute for a row that already matched an earlier candidate. The finding below reproduces this new query failure.
  • Compatibility analysis: Compared Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Their relevant conditional and nonempty IN semantics agree. Spark evaluates later membership candidates only until that row matches. The fix needs to preserve that behavior.
  • Key design decisions: Placing the empty-batch guard before eager/lazy selection protects both paths. The separate column check protects nested membership even when another expression returns a misleading scalar.
  • Implementation sketch: Two small production changes reuse new_empty_array and collect_columns. Rust tests cover empty results and candidate classification. The SQL fixture exercises conditional candidates, negation, structs, and arrays.
  • Performance: The added column traversal occurs during expression construction. Empty-array allocation is confined to empty batches. Literal-only CASE expressions lose native scalar classification, although Spark normally folds them first. No additional evidence-backed P1/P2 performance issue was identified. No benchmark was run.
  • Design: The producer and consumer checks are focused and understandable. The newly exercised dynamic membership path still needs row-level short-circuiting before this can safely handle fallible candidates.
  • Abstraction & complexity: No new production abstraction is introduced. Reusing Arrow allocation and DataFusion column discovery is appropriate. The extended test Probe directly exercises the defensive check.
  • Behavioral changes worth calling out: Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, whose conditional path delegates to DataFusion CaseExpr, the intended change fixes erroneous NULL caching. A native comparison also demonstrates the unintended transition from a successful result to an unreachable-cast error. Configuration, serialization, and support levels are unchanged.
  • Suggested improvements: Preserve per-row short-circuiting for dynamic IN, or fall back for affected fallible candidates. Add the mixed-match ANSI regression described below. No other issue meeting the P1/P2 reporting bar was identified.

Exact-head CI: Comet CI passed the native build, Rust tests, Spark 4.1 Comet suites, and TPC-H/TPC-DS checks. Logs explicitly confirm the new Rust tests and in_case_when_candidate.sql passed. Spark’s own SQL suites, Iceberg suites, macOS, and benchmarks were skipped.

Validation: All 30 focused existing Rust tests passed. Disposable native probes demonstrated both the intended fix and the reported regression. Spark 4.1.3 directly executed the finding’s SQL and returned the expected rows. The native reproduction used the head expressions and the base-equivalent lazy CaseExpr path. A full Comet JVM reproduction and full base/release builds were not run. Runtime testing of other Spark versions was not performed. Disposable project tests were removed and the checkout is clean.

// as IN does to build its list once, so return an empty array instead.
if batch.num_rows() == 0 {
let data_type = self.data_type(&batch.schema())?;
return Ok(ColumnarValue::Array(new_empty_array(&data_type)));

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] Preserve per-row IN short-circuiting when making CASE candidates dynamic. With ANSI enabled and both rows in one native batch, SELECT id, id IN (0L, CASE WHEN id = 0 THEN CAST('bad' AS BIGINT) WHEN id = 2 THEN id END) FROM range(0, 2, 1, 1) should return (0,true), (1,NULL). This empty-array return routes the candidate into DataFusion's dynamic IN, which evaluates later candidates against the entire batch unless every row already matched. It therefore executes the invalid cast for id=0, despite that row having matched 0L, and aborts the query. The underlying dynamic loop predates this PR, but this change exposes it for a CASE shape whose base lazy path returned the expected result. Please evaluate later candidates only for unmatched rows, or route affected fallible expressions through Spark, and add this mixed-match regression.

Evidence: Spark 4.1.3 with spark.sql.ansi.enabled=true executed the SQL and returned [Row(id=0, matched=True), Row(id=1, matched=None)]; its optimized plan retained both CASE branches. A disposable Rust reproduction used an Int64 batch [0,1], in_list, the same two CASE branches, and Comet's ANSI Cast from 'bad' to Int64. The base-equivalent lazy CaseExpr path returned [true,NULL]. Head CaseWhenExpr returned Err(External(CastInvalidValue { value: "bad", from_type: "STRING", to_type: "BIGINT" })). Base source selects that lazy path because this cast is fallible. DataFusion 55.1.0 InListExpr::evaluate calls each remaining candidate on the original batch and only short-circuits when the entire accumulated result is true. Reproduction source was retained at /tmp/comet-review-6763-repro.rs; no project changes remain.

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.

I reproduced this on Spark 4.1 with ANSI on, and it does not need a CASE. SELECT id, id IN (0L, CAST(s AS BIGINT)) FROM t raises CAST_INVALID_INPUT when s is 'bad' on the row where id = 0, and it does the same on main before this PR. InListExpr and NestedPredicate evaluate every later candidate over the whole batch and only stop once every row has matched, where Spark stops per row. So this PR does not add the behavior. A CASE candidate now goes through the same loop instead of being read as a constant NULL. That happened to be right for this shape and wrong for others. id IN (0L, CASE WHEN id = 0 THEN CAST(s AS BIGINT) WHEN id = 1 THEN 1L END) gave NULL for id = 1 where Spark returns true.

Could you open an issue for the whole-batch evaluation of later IN candidates, link it from #6006, and reference it here? I'd rather track it than grow this PR into a per-row IN evaluator, but Spark 4 runs with ANSI on by default, so I don't want it lost.

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.

Thanks both. As @andygrove found, this predates the PR: InListExpr and NestedPredicate evaluate every later candidate over the whole batch, and id IN (0L, CAST(s AS BIGINT)) already fails this way on main. This PR moves a column-reading CASE from the frozen constant onto that same path. I opened #6823 for the whole-batch evaluation of later candidates and linked it from #6006, and left the per-row evaluation out of this PR.

-- no ELSE
query
SELECT id, id IN (CASE WHEN id <> 1 THEN id END), id NOT IN (CASE WHEN id <> 1 THEN id END)
FROM range(0, 3)

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.

CometCoalesce builds the same CaseWhen proto, so coalesce hit this bug too, and the fixture has no query for it. On main before this PR, id IN (coalesce(v, id)) over a nullable column v returns NULL for every row. Could we add SELECT id, id IN (coalesce(nullif(id, 1), 5L)) FROM range(0, 3) and a CASE with two WHEN branches and no ELSE, such as SELECT id, id IN (CASE WHEN id = 0 THEN id WHEN id = 2 THEN id END) FROM range(0, 3)? The first fails without the fix, and both pass with it.

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 both, and the spark_partition_id() query @comphead suggested. Each of the three fails without the fix.

if is_volatile(child) {
// A candidate that reads a column is not a constant, even if it returns a scalar
// for the empty batch
if is_volatile(child) || !collect_columns(child).is_empty() {

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.

The new collect_columns check only protects nested operands. Any other type returns into DataFusion's in_list at the top of spark_in_list, and InListExpr::try_new decides there by evaluating the candidates on an empty batch. A column-reading expression that returns a scalar for zero rows would bring this bug back for flat IN, with no test to notice. Could we say that in a comment on that early return, and add a sentence to adding_a_new_expression.md saying an expression that reads a column must return an array for an empty batch?

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 a comment on that early return, pointing at apache/datafusion#26082, and a "Returning a scalar or an array" section to adding_a_new_expression.md.

@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 @viirya. The early return and the spark_in_list check fix the frozen NULL for CASE and IF candidates that read a column, and I did not repeat the points already in the thread. The inline comments ask about a branch-1.1 backport, suggest a fixture query for a CASE that Spark cannot fold, point at two inert config lines, a test row that cannot fail without the fix and tests the SQL file already covers, and link the upstream DataFusion fix for the same bug.

Comment on lines +353 to +355
// No row chooses a branch, which both the eager and the lazy evaluation answer with a
// scalar NULL. A scalar from an empty batch is taken to mean the expression is constant,
// as IN does to build its list once, so return an empty array instead.

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.

Should this also go to branch-1.1? There CASE, IF and nullif evaluate through DataFusion's CaseExpr, which treats a zero-row mask as all true, so IF(id = 1, NULL, id) returns its NULL literal as a scalar for the empty batch. From reading the code I expect id IN (IF(id = 1, NULL, id)) to return NULL for every row in 1.1.0 as well. I have not run it. CaseWhenExpr is not on that branch, so a backport would need this check in that branch's IfExpr::evaluate and in a wrapper around the CaseExpr that create_case_expr returns. If that holds, it is a backport-1.1 candidate under backporting.md.

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.

Yes, I think it should. branch-1.1 uses the same DataFusion 55.1.0, and its IfExpr::evaluate delegates to CaseExpr, which returns a NULL THEN literal as a scalar for an empty batch, the same lazy path I confirmed returns a scalar NULL on main. I'll open a branch-1.1 backport after this merges, with the check in IfExpr::evaluate and in a wrapper around the CaseExpr from create_case_expr, and verify the reproduction there first.

Comment on lines +18 to +20
-- IN builds its list once when every candidate is a constant, which it decides by evaluating the
-- candidates on an empty batch and checking for a scalar. A CASE or IF that depends on a column
-- must not return a scalar NULL there, or IN compares every row against NULL.

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.

The description says a CASE made only of literals just loses the static list, and that Spark folds such a CASE first. Spark cannot fold one over spark_partition_id(), but Comet plans that function as a Literal (SparkPartitionIdBuilder). On main the eager path answers the empty-batch probe with a scalar NULL, so from reading the code I expect SELECT id, id IN (CASE WHEN spark_partition_id() = 0 THEN 1L END, 5L) FROM range(0, 3, 1, 1) to return NULL on every row, where Spark returns false, true, false. Could we add it here? It would also catch a later change that skips the early return for a CASE without columns. I have not run it.

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. Without the fix it returns NULL on every row, as you expected.

Comment on lines +23 to +24
-- Config: spark.comet.sparkToColumnar.enabled=true
-- Config: spark.comet.sparkToColumnar.supportedOperatorList=Range

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.

With spark.comet.exec.range.enabled=true, CometExecRule turns range(...) into CometRangeExec and only falls back to the Spark-to-Arrow conversion when that operator declines a range (around line 504 of CometExecRule.scala). So the two sparkToColumnar lines have no effect here, and naming Range in spark.comet.sparkToColumnar.supportedOperatorList is deprecated per its doc in CometConf.scala. Could we drop both lines? The file would then fail if native Range ever stopped taking these ranges, rather than quietly running on a converted leaf.

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.

Dropped both lines.

Comment on lines +1268 to +1269
// A branch that can fail is evaluated lazily
(vec![(a_is_1(), a_div_b)], None, false),

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.

From reading DataFusion 55.1.0, this row passes without the early return. For CASE WHEN a = 1 THEN a / b END, CaseExpr picks expr_or_expr, treats the zero-row mask as all true and returns a / b evaluated on the empty batch, which is already an empty array. Could the lazy row be (vec![(a_is_1(), null())], Some(a_div_b), false) instead? CaseExpr returns its NULL literal as a scalar there, so the row would fail if the check moved into evaluate_eagerly. The three eager rows reach the same early return before the path is chosen, so one of them would do. I have not run it.

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.

You're right, that row already returned an empty array without the fix. It is now THEN NULL ELSE a / b, which returns a scalar NULL without it, and I kept one eager row.

Comment on lines +1284 to +1286
/// IN takes a candidate that returns a scalar for an empty batch as a constant.
#[test]
fn in_list_does_not_take_a_case_as_constant() {

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.

This test builds the same two candidates as the first and third queries of in_case_when_candidate.sql and expects the same [true, NULL, true], which the SQL file already checks against Spark through the same in_list. Could we drop it? The case_b candidate in candidates_reading_a_column_remain_dynamic could go too. Its list type fails can_merge, so it is evaluated lazily, and DataFusion's InfallibleExprOrNull path already returns an empty array for it, so from reading the code it passes on main without either change. The Probe candidate is the one that tests the new check. I have not run it.

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.

Confirmed both, and removed in_list_does_not_take_a_case_as_constant and the case_b candidate. Without the fix only the Probe candidate failed.

Comment on lines +380 to +382
// A candidate that reads a column is not a constant, even if it returns a scalar
// for the empty batch
if is_volatile(child) || !collect_columns(child).is_empty() {

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.

DataFusion has the same bug in InListExpr::try_new, filed as apache/datafusion#26082, and the approved apache/datafusion#26083 fixes it there by skipping the empty-batch probe for an item with a leaf that is not a Literal. Could the description link it, so the next DataFusion upgrade knows the flat path gets a similar guard? This check could also be one walk, as calls_jvm in parquet_exec.rs does with exists, for example child.exists(|e| Ok(e.is_volatile_node() || e.is::<Column>())).unwrap_or(true). collect_columns visits every node and clones each Column into a HashSet just to test that it is empty, after is_volatile has already walked the tree.

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.

Linked both in the description. The check is now a single exists walk with the same rule as apache/datafusion#26083: a candidate is a constant only when every leaf is a Literal and no node is volatile.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:expressions Expression evaluation bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

IN treats a column-dependent CASE WHEN / IF candidate as the constant NULL

4 participants