Repository navigation
Conversation
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.
sunchao
left a comment
There was a problem hiding this comment.
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, ornullifcandidate into a cached scalar NULL. Queries such asid 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_listconstant detection. - Correctness: The intended membership correction is covered by passing tests. However, making these candidates dynamic exposes whole-batch
INevaluation: 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
INsemantics 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_arrayandcollect_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
Probedirectly exercises the defensive check. - Behavioral changes worth calling out: Compared with
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, whose conditional path delegates to DataFusionCaseExpr, 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))); |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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() { |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
| // 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. |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| -- 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. |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Added. Without the fix it returns NULL on every row, as you expected.
| -- Config: spark.comet.sparkToColumnar.enabled=true | ||
| -- Config: spark.comet.sparkToColumnar.supportedOperatorList=Range |
There was a problem hiding this comment.
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.
| // A branch that can fail is evaluated lazily | ||
| (vec![(a_is_1(), a_div_b)], None, false), |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| /// 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() { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| // 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() { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
Which issue does this PR close?
Closes #6762.
Rationale for this change
INbuilds its candidate list once when every candidate is a constant. Both DataFusion'sInListExprand Comet'sspark_in_listdecide that by evaluating each candidate on an empty batch and checking for a scalar.CaseWhenExpr, whichIF,nullifandcoalescealso 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 exampleid IN (IF(id = 1, NULL, id))returned NULL for every row.DataFusion has the same bug in
InListExpr::try_newfor its ownCaseExpr(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 plansspark_partition_id()as a literal, soCASE WHEN spark_partition_id() = 0 THEN 1L ENDstill 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::evaluatereturns an empty array of the result type for an empty batch. The check runs before both the eager evaluation from perf: optimize nativeCASE WHENandIF(up to 12x faster) #6350 and the lazy one, because DataFusion'sCaseExpralso 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_listtreats a candidate as a constant only when every leaf is aLiteraland 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.mdsays 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
PhysicalExprimplementations and scalar functions return a scalar only when all their inputs are scalars, or when a NULL scalar argument makes every row NULL.INstill 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?
spark_in_listkeeps 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.sqlruns the queries from the issue with nativerange: IF,nullif,coalesce, CASE with and without ELSE,NOT IN, a CASE overspark_partition_id(), and struct and array operands, including a double leaf so thatspark_in_listis used. It fails without the fix.This pull request and its description were written by Isaac.