Skip to content

fix: dispatch StaticInvoke and Invoke only into Spark's own classes - #6542

Merged
andygrove merged 4 commits into
apache:mainfrom
andygrove:fix/spark-only-invoke-dispatch-6425
Oct 2, 2026
Merged

andygrove merged 4 commits into
apache:mainfrom
andygrove:fix/spark-only-invoke-dispatch-6425

Conversation

@andygrove

@andygrove andygrove commented Oct 2, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6425.

This replaces #6455, which tried to fix the same issue inside the dispatcher.

Rationale for this change

Since #5692, CometStaticInvoke and CometInvoke send every call they don't otherwise handle to the codegen dispatcher. That includes calls to DataSource V2 catalog functions, which are user code. Such a function can return a Decimal at a different scale from the type it declares, or one that doesn't fit it. Spark only corrects that when it writes a row, and any expression around the call reads the value as returned. The dispatcher has to write an Arrow vector of the declared type, so whatever reads its output sees a different value.

#6455 rescaled the value in the dispatcher and tried to keep every consumer of the call in the same kernel or in Spark. But where Spark writes rows depends on whole-stage codegen, and each review round found another path: transitive parents, projected aliases, explode, a dispatched map(...) around the call, and ApplyFunctionExpression. In 1.0.0 these calls fell back to Spark, so this PR goes back to that. The dispatcher's handling of StaticInvoke and Invoke becomes an allow list instead of a catch-all.

What changes are included in this PR?

The new CometInvokeTargets decides which calls the dispatcher may run. A StaticInvoke or Invoke is allowed when the class it calls is part of Spark (org.apache.spark.sql.catalyst.* or org.apache.spark.unsafe.*) and isn't a ScalarFunction. An Invoke on a Catalyst value, such as Spark 4's is_valid_utf8 calling isValid on a UTF8String, is allowed too. An ApplyFunctionExpression never is. The one exception is the predicate of a typed Dataset.filter. Spark calls it through scala.Function1 or FilterFunction, and it returns a boolean, so there's nothing for a row write to correct. #6497 covers that path, so it stays on the dispatcher.

emitJvmCodegenDispatch applies the check to the whole tree rather than just its root, because the kernel runs every node under the dispatched expression. That's what catches a DSv2 call nested inside map(...) or another dispatched expression.

Spark's own lowerings keep dispatching, for example binary lpad/rpad, encode/to_binary, to_time and the Spark 4 evaluator Invokes. I checked the StaticInvoke and Invoke targets in Spark 3.4.3 through 4.2.0, and every function lowering calls into one of those two packages. Iceberg functions with a native handler are unaffected. An Iceberg function with no handler, or nested inside a dispatched expression, now falls back. So does a typed filter whose input decoding calls a class outside Spark, for example a Java enum's valueOf. The dispatcher's decimal writer is unchanged. The user guide now says that other DSv2 functions run in Spark.

How are these changes tested?

CometCodegenSuite registers a DSv2 function catalog whose functions return scale-0 decimals for DECIMAL(10, 2), one for each lowering (Invoke, StaticInvoke and ApplyFunctionExpression). It runs the shapes from the #6455 reviews:

  • the plain result, IS NULL and a cast to string
  • abs(...), and array and struct access on the result
  • count, max and sum
  • a projected alias, explode, and map(...) around the call
  • DISTRIBUTE BY on the call

Each one falls back with the new reason and matches Spark. A unit test pins which calls the allow list accepts and declines, both at the root and nested under map(...). CometIcebergSystemFunctionSuite adds @sunchao's map_values(map('k', truncate(10, d)))[0] case.

The #5573 and #5575 tests used test classes as Invoke targets, which the dispatcher now declines. The #5573 test now uses a UDF that captures a non-serializable object, and the #5575 test uses an Invoke on a UTF8String. The unlisted StaticInvoke test in the Iceberg suite now expects a fallback.

With the check disabled, both new SQL tests fail. The DSv2 one returns #6425's wrong answers (0.03 for 3.00, 1000000.00 where Spark returns null).

Local runs:

  • Spark 4.1: CometCodegenSuite and CometIcebergSystemFunctionSuite (121 tests), CometExpressionSuite (174), and CometSqlFileTestSuite, CometStringExpressionSuite, CometCodegenSourceSuite, CometCodegenFuzzSuite and the two Iceberg function extension and pushdown suites (730).
  • Spark 3.4, 3.5, 4.0 and 4.2: CometCodegenSuite, CometIcebergSystemFunctionSuite and CometStringExpressionSuite (158, 158, 160 and 146). The cancellations are version guards, and there's no Iceberg build for Spark 4.2.

Since apache#5692, CometStaticInvoke and CometInvoke send every call they do not
otherwise handle to the codegen dispatcher, including calls to DataSource V2
catalog functions. Such a function can return a Decimal at a different scale
from its declared type, or one that does not fit it. Spark corrects that only
when it writes a row, and an expression around the call reads the value as
returned, so the dispatcher's Arrow output cannot match it (apache#6425).

The dispatcher now runs a StaticInvoke or Invoke only when it calls one of
Spark's own classes, or the predicate of a typed Dataset.filter, and never an
ApplyFunctionExpression. It checks the whole tree it is given, since the kernel
runs every node in it. Anything else falls back to Spark, as in 1.0.0.

Closes apache#6425.
@andygrove andygrove added bug Something isn't working regression A bug that did not affect the most recent Comet release area:expressions Expression evaluation backport-1.1 Candidate for backporting to 1.1 release branch run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue labels Oct 2, 2026
The decline reason no longer claims that only Spark's own code runs in the
dispatcher, which the typed filter predicate contradicts. The DSv2 fixture no
longer forces a single file, which only apache#6455's writer tests needed. The
dispatch benchmark's description and the Iceberg guide match the allow list.
On main's writer the partition ids happen to match Spark for these values, so
comparing answers alone could not catch the call being dispatched.

@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: Dispatching arbitrary DSv2 calls could change decimal scale and overflow behavior when their results crossed into Arrow.
  • Design approach: CometInvokeTargets checks the entire dispatched expression tree and declines external invocation targets and DSv2 functions.
  • Correctness / compatibility analysis: The guard matches Spark’s three DSv2 lowerings. Source comparisons covered Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No introduced P1/P2 issues found within this review.
  • Key design decisions: Spark helpers, Catalyst-value receivers and boolean typed-filter predicates remain eligible. Existing native Iceberg handlers remain available.
  • Implementation sketch: One shared check runs before binding and serialization. This centralizes the policy and adds a planning-time tree traversal without adding per-row work.
  • Behavioral changes worth calling out: External calls and Iceberg calls nested inside dispatched expressions now cause Spark fallback. The planner keeps dependent consumers in Spark until an existing execution boundary.
  • Suggested improvements: None at P1/P2 priority.

Reviewed full SHA 179fdd33dac82c33066d6195ed1812f5bb756b86 against base a2d7370346810fb43f7c24fd96042eb31acd6025, covering all three PR commits and all seven PR-changed files. The additional array-kernel differences between the endpoint trees come from base-only #6518, not PR deletions. The PR is not a draft. Snapshot and live discussions contained no existing review concerns.

Routed skills: review-comet-pr and review-comet-expression-pr.

Exact-head CI: 17 checks passed, 13 skipped, and three remained in progress. No failures were reported. Spark compilation profiles and strict Scala warnings passed. Rust tests, the native build and Spark 4.1 SQL-test preparation remained running. Iceberg integration jobs were skipped.

Validation: A standalone harness passed 40 checks using the exact-head guard and Spark 4.1.3, covering permitted calls, all three DSv2 forms, nesting and decimal materialization. Only the diagnostic-name helper was stubbed. Full integration suites could not start because cached Maven dependencies disappeared during setup and Maven Central access was blocked. No full native rebuild or local Spark/Iceberg integration verdict was obtained.

Fold the typed filter rule into the callee lookup and name the DSv2 check
once. The tests use Spark's InMemoryCatalog and the suite's withTypedCol
instead of their own catalog and table setup, derive the static-invoke
fixture from the shared one, drop a leftover Invoke target, and pin how each
fixture lowers.
@andygrove andygrove mentioned this pull request Oct 2, 2026
20 of 21 tasks
@andygrove
andygrove added this pull request to the merge queue Oct 2, 2026
@andygrove

Copy link
Copy Markdown
Member Author

Thanks @sunchao This is one of the last blocking issues for 1.1.0

Merged via the queue into apache:main with commit b0f0de3 Oct 2, 2026
50 checks passed
andygrove added a commit that referenced this pull request Oct 2, 2026
…6542) (#6554)

* fix: dispatch StaticInvoke and Invoke only into Spark's own classes

Since #5692, CometStaticInvoke and CometInvoke send every call they do not
otherwise handle to the codegen dispatcher, including calls to DataSource V2
catalog functions. Such a function can return a Decimal at a different scale
from its declared type, or one that does not fit it. Spark corrects that only
when it writes a row, and an expression around the call reads the value as
returned, so the dispatcher's Arrow output cannot match it (#6425).

The dispatcher now runs a StaticInvoke or Invoke only when it calls one of
Spark's own classes, or the predicate of a typed Dataset.filter, and never an
ApplyFunctionExpression. It checks the whole tree it is given, since the kernel
runs every node in it. Anything else falls back to Spark, as in 1.0.0.

Closes #6425.

* test: tidy the #6425 fallback reason, fixture and docs

The decline reason no longer claims that only Spark's own code runs in the
dispatcher, which the typed filter predicate contradicts. The DSv2 fixture no
longer forces a single file, which only #6455's writer tests needed. The
dispatch benchmark's description and the Iceberg guide match the allow list.

* test: assert that a DSv2 call in a shuffle's partitioning falls back

On main's writer the partition ids happen to match Spark for these values, so
comparing answers alone could not catch the call being dispatched.

* refactor: simplify the dispatcher allow list and its #6425 tests

Fold the typed filter rule into the callee lookup and name the DSv2 check
once. The tests use Spark's InMemoryCatalog and the suite's withTypedCol
instead of their own catalog and table setup, derive the static-invoke
fixture from the shared one, drop a leftover Invoke target, and pin how each
fixture lowers.

(cherry picked from commit b0f0de3)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:expressions Expression evaluation area:Iceberg 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 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.

Decimal-returning DSv2 functions lose their scale through the codegen dispatcher

2 participants