Skip to content

fix: defer throwing literal casts to Spark at runtime - #6035

Open
sunchao wants to merge 3 commits into
apache:mainfrom
sunchao:upstream/fix-literal-cast-planning
Open

sunchao wants to merge 3 commits into
apache:mainfrom
sunchao:upstream/fix-literal-cast-planning

Conversation

@sunchao

@sunchao sunchao commented Sep 19, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

Preserve Spark's conditional evaluation when a literal cast throws during Comet planning. Spark can deliberately leave an invalid cast unfolded inside a conditional branch; planning must not throw for a branch that is never reached.

The failed literal cast keeps the containing projection in Spark. This deliberately avoids batch codegen dispatch, which can evaluate a row that Spark skips after a LIMIT. A shared explanatory fallback reason and comments document that decision. Successful literal casts and non-ANSI casts that return null still retain native projection.

How was this patch tested?

Spark 4.1 reactor compilation, formatting/style checks, and two focused CometNativeCastSuite tests passed against the matching native library. Coverage asserts the invalid literal cast survives optimization, checks Spark error parity through the shared helper, exercises LIMIT with dispatch enabled and disabled, and verifies native controls for successful casts and non-ANSI invalid casts. The shared error helper retains native-plan checking by default and allows this intentional fallback test to disable only that check.

Rebased on main ac9ae94. Hosted current-head checks are pending.

@sunchao sunchao added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 19, 2026
@github-actions github-actions Bot added bug Something isn't working area:expressions Expression evaluation labels Sep 19, 2026

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

Thanks for tracking this down — the root cause analysis is accurate and I verified it against Spark's source. ConstantFolding.tryFold deliberately leaves a throwing expression in place when it sits in a conditional branch, tagging it FAILED_TO_EVALUATE rather than folding it:

// sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/expressions.scala
case NonFatal(_) if isConditionalBranch =>
  expr.setTagValue(FAILED_TO_EVALUATE, ())
  expr

So Cast(Literal("bad"), IntegerType) survives into the physical plan, and Comet's unconditional cast.eval() during serialization raises for a branch that may never be visited at runtime. That turns a query Spark answers successfully into a planning-time failure. Real bug, worth fixing, and catching NonFatal matches the range Spark's own tryFold catches, so the two stay aligned on what counts as a recoverable failure.

My main concern is where the fallback happens rather than whether it should.

The fallback bypasses the codegen dispatcher. CometCast mixes in CodegenDispatchFallback, whose contract states that Unsupported means "no native path exists for this case; run Spark's doGenCode inside the Comet pipeline," and that Spark fallback is reserved for cases the dispatcher itself cannot handle. But exprToProtoInternal only calls dispatchIfFallback from the Unsupported and Incompatible branches — Compatible goes straight to convert, so a None returned there falls the whole projection back to Spark without the dispatcher ever being tried.

A throwing literal cast is squarely within what the dispatcher handles: Cast.doGenCode compiles, and it raises at runtime only if the branch is actually visited, which is exactly Spark's behavior. Evaluating in getSupportLevel and returning Unsupported(Some(reason)) on failure would let the framework try the dispatcher first and fall back to Spark only if that fails — strictly better than the current outcome, and it is what this same file already does for VariantType:

// Reporting `Unsupported` lets the `CodegenDispatchFallback` mixin try the dispatcher
// and then fall back to Spark cleanly.

It also removes the inconsistency of getSupportLevel reporting Compatible for an expression convert cannot actually produce. The cost is evaluating the literal twice on the success path, which seems negligible — the failure path never reaches convert.

Worth noting for the record: ConstantFolding.FAILED_TO_EVALUATE would be a more precise signal than catching from eval(), but it is private[sql] and unreachable from org.apache.comet.expressions, so try/catch is a reasonable substitute. A comment saying so would save the next reader from re-deriving it.

I checked for sibling instances of this bug class and found none that need fixing here. I swept every eager .eval( call in the serde layer on main. Round's r.scale.eval(EmptyRow) looked like the closest match, but RoundBase.dataType and checkInputDataTypes already evaluate _scale during analysis, so Spark raises first; the rest (ArraySort's ascendingOrder, the percentile and bloom-filter parameters, window frame bounds) are all analysis-validated foldable arguments. The scope of this PR looks right.

Smaller points, all in inline comments: the fallback reason should be a shared constant to match the convention this file documents, the test duplicates the existing checkSparkError helper in a slightly weaker form, and it is in a suite about something else. A negative test asserting that successful literal casts still fold natively would guard against the catch being widened later.

Comment thread spark/src/main/scala/org/apache/comet/expressions/CometCast.scala
Comment thread spark/src/test/scala/org/apache/comet/serde/CometScalarFunctionSuite.scala Outdated
Comment thread spark/src/test/scala/org/apache/comet/serde/CometScalarFunctionSuite.scala Outdated
Comment thread spark/src/test/scala/org/apache/comet/serde/CometScalarFunctionSuite.scala Outdated
Comment thread spark/src/test/scala/org/apache/comet/serde/CometScalarFunctionSuite.scala Outdated

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

I agree with @viirya that reporting Unsupported from getSupportLevel fits the CodegenDispatchFallback contract, but I don't think it's a drop-in change for this test. With the dispatcher, the projection stays native and evaluates the whole batch. Both rows of cast_branch_rows land in one batch, so row 1 takes the id = 1 branch and the dispatched CAST('bad' AS INT) throws, even though Spark stops after row 0 because of the LIMIT 1. Falling the projection back to Spark, as this PR does now, is what keeps Spark's answer for that query.

So I think we need to pick one on purpose. If we go with the dispatcher, could the masked query use a predicate no row satisfies, such as id = 5, so it tests the literal cast rather than how a native projection behaves under a limit? If we keep the Spark fallback, could CometCast.convert get a comment saying it's deliberate, so nobody later moves it into getSupportLevel and breaks the LIMIT case?

@sunchao

sunchao commented Sep 23, 2026

Copy link
Copy Markdown
Member Author

Thanks! Agree it is not ideal to fallback to Spark fully for this scenario. I should have checked before Codex-ing this PR lol. Let me revise it and then ask for another round of reviews.

@sunchao sunchao left a comment

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.

Summary

  • Prior state and problem: Eager cast.eval() could fail planning for an invalid literal in a conditional branch Spark never evaluates.
  • Design approach: Catch NonFatal, attach a fallback reason, and defer evaluation to Spark. Successful literal casts retain their existing folding path.
  • Correctness / compatibility analysis: Spark’s conditional-folding rules support this behavior across the checked versions. The regression verifies both the successful LIMIT 1 query and matching CAST_INVALID_INPUT exceptions when the bad row is demanded. No introduced P1/P2 issues found within this review.
  • Key design decisions: Returning None bypasses codegen dispatch and preserves Spark’s demand-driven evaluation. The existing dispatcher discussion identifies a real tradeoff: processing a full batch can encounter the bad second row before the limit stops execution. I found no substantiated unresolved P1/P2 concern in the existing feedback.
  • Implementation sketch: A localized exception handler uses the existing fallback mechanism without adding an abstraction, duplicate successful evaluations, or per-row native work.
  • Behavioral changes worth calling out: A previously failing query can complete through a Spark projection while retaining the native Parquet scan. Demanded invalid casts still throw.
  • Suggested improvements: None additional at P1/P2 severity.

Reviewed the entire two-file diff from 5fdc96199685061b7c67ad28651b4c0a3dcd6541 to 917b0510acc056f2d8e15f10ccc7c3548f6c59cf. The PR is not a draft. Read all existing reviews, issue comments, inline comments, and threads. Routed skills: review-comet-pr and audit-comet-expression, scoped to the changed literal-cast behavior.

Exact-head CI: All four Comet test groups and all seven Spark 4.1 SQL matrix jobs passed. The expressions log confirms the new regression passed, with 1,513 tests successful overall. CI used merge commit 629b5821151c5f76a46368c0453de040ed673260, whose tree matches the reviewed head. The required gate remains red because the Spark 3.5 lint job failed resolving the scalafix Maven plugin, before source compilation. This does not establish a PR-introduced defect.

Validation limits: No local build or runtime tests were rerun because this checkout lacks native/JVM build artifacts and Spark dependencies. Validation used the exact-tree CI logs and relevant Spark source/tests for 3.4.3, 3.5.8/3.5.9, 4.0.1/4.0.4, 4.1.1/4.1.3, and 4.2.0. Runtime coverage was verified for Spark 4.1 only. No project code or GitHub state was changed.

@sunchao
sunchao force-pushed the upstream/fix-literal-cast-planning branch from 917b051 to cae771a Compare October 4, 2026 17:50

@sunchao sunchao left a comment

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.

Summary

  • Prior state and problem: Comet eagerly evaluated literal casts during planning, including invalid casts Spark deliberately preserves inside conditional branches. This could fail a query before its LIMIT avoided the offending row.
  • Design approach: Catch NonFatal from cast.eval(), attach a shared explanatory fallback reason, and leave the containing operator in Spark.
  • Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Spark source and optimizer tests confirm the relevant semantics across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Demanded invalid casts still raise CAST_INVALID_INPUT. Successful casts and non-ANSI invalid casts retain native projection.
  • Key design decisions: Keeping the fallback in convert deliberately avoids batch codegen dispatch, which could evaluate a row Spark never requests. The existing dispatcher discussion is addressed by this documented choice. No substantiated existing P1/P2 concern remains unresolved.
  • Implementation sketch: The change uses the existing Option and fallback-tag mechanisms without adding an abstraction or duplicating successful evaluation. The shared error helper retains native-plan checking by default. No new per-row work is introduced.
  • Behavioral changes worth calling out: Compared with branch-1.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, the affected query intentionally changes from a planning failure to Spark runtime evaluation, retaining its native scan. Other differences in the touched files were already present in the requested base.
  • Suggested improvements: None additional at P1/P2 severity.

Reviewed the entire three-file, three-commit diff from ac9ae94d057013095fdf1b70e2d4bc5b24a58d90 to cae771a9db078dac28473841ceb7e15c44e02990. The PR remains open and non-draft. Read the supplied reviews, issue comments, inline comments and review threads. Routed skills: review-comet-pr and review-comet-expression-pr.

Exact-head CI: 36 checks passed, 15 skipped, none failed or pending. This includes Required Checks, all four Comet test groups and all nine Spark 4.1 SQL jobs. The expressions log confirms both new regressions passed, with 2,065 successful tests overall. CI tested merge commit 54262dafafee71d3a4b50f27377af548eff81cc3, whose tree matches the reviewed head.

Validation limits: A local Spark 3.5.9 reproduction confirmed the retained invalid cast, eager-evaluation failure, successful LIMIT 1, demanded-row error and ANSI/non-ANSI controls. This was a Spark reference test, not a local Comet run. Comet runtime validation relied on exact-tree Spark 4.1.3 CI. No local Comet build or runtime matrix for other Spark versions was run. Project code and GitHub state were unchanged.

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 run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants