fix: dispatch StaticInvoke and Invoke only into Spark's own classes - #6542
Conversation
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.
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
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Dispatching arbitrary DSv2 calls could change decimal scale and overflow behavior when their results crossed into Arrow.
- Design approach:
CometInvokeTargetschecks 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.
|
Thanks @sunchao This is one of the last blocking issues for 1.1.0 |
…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)
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,
CometStaticInvokeandCometInvokesend 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 aDecimalat 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 dispatchedmap(...)around the call, andApplyFunctionExpression. In 1.0.0 these calls fell back to Spark, so this PR goes back to that. The dispatcher's handling ofStaticInvokeandInvokebecomes an allow list instead of a catch-all.What changes are included in this PR?
The new
CometInvokeTargetsdecides which calls the dispatcher may run. AStaticInvokeorInvokeis allowed when the class it calls is part of Spark (org.apache.spark.sql.catalyst.*ororg.apache.spark.unsafe.*) and isn't aScalarFunction. AnInvokeon a Catalyst value, such as Spark 4'sis_valid_utf8callingisValidon aUTF8String, is allowed too. AnApplyFunctionExpressionnever is. The one exception is the predicate of a typedDataset.filter. Spark calls it throughscala.Function1orFilterFunction, 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.emitJvmCodegenDispatchapplies 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 insidemap(...)or another dispatched expression.Spark's own lowerings keep dispatching, for example binary
lpad/rpad,encode/to_binary,to_timeand the Spark 4 evaluatorInvokes. I checked theStaticInvokeandInvoketargets 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'svalueOf. The dispatcher's decimal writer is unchanged. The user guide now says that other DSv2 functions run in Spark.How are these changes tested?
CometCodegenSuiteregisters a DSv2 function catalog whose functions return scale-0 decimals forDECIMAL(10, 2), one for each lowering (Invoke,StaticInvokeandApplyFunctionExpression). It runs the shapes from the #6455 reviews:IS NULLand a cast to stringabs(...), and array and struct access on the resultcount,maxandsumexplode, andmap(...)around the callDISTRIBUTE BYon the callEach 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(...).CometIcebergSystemFunctionSuiteadds @sunchao'smap_values(map('k', truncate(10, d)))[0]case.The #5573 and #5575 tests used test classes as
Invoketargets, which the dispatcher now declines. The #5573 test now uses a UDF that captures a non-serializable object, and the #5575 test uses anInvokeon aUTF8String. The unlistedStaticInvoketest 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.03for3.00,1000000.00where Spark returns null).Local runs:
CometCodegenSuiteandCometIcebergSystemFunctionSuite(121 tests),CometExpressionSuite(174), andCometSqlFileTestSuite,CometStringExpressionSuite,CometCodegenSourceSuite,CometCodegenFuzzSuiteand the two Iceberg function extension and pushdown suites (730).CometCodegenSuite,CometIcebergSystemFunctionSuiteandCometStringExpressionSuite(158, 158, 160 and 146). The cancellations are version guards, and there's no Iceberg build for Spark 4.2.