Repository navigation
Conversation
… evaluation Registering a native UDF with a type Comet has no native representation for, such as a user-defined type, succeeded, and every call then fell back to Spark, which cannot evaluate it. CometNativeUDF.register now refuses such a signature before it loads the library. Spark also evaluates a call itself while planning a query: over local data, in a filter on partition columns, and to sample the keys of a global sort. The error a native UDF raises there said that Comet had not taken the operator holding the call, which sent users to the explain output for a fallback that never happened. The error, the scaladoc and the Rust UDF guide now list those cases, and the guide says to sort on a column holding a UDF's result rather than on the call. Split out of apache#6697.
sunchao
left a comment
There was a problem hiding this comment.
No introduced P1/P2 issues found within this review.
Reviewed the full five-file diff from base 7d294535e87367d819c4e2a23422cd1158ac31b8 to head cd17895230f896b19877e25dbed531982157fcd2. The PR is not a draft. The snapshot and live discussion contained no reviews or comments, so there were no existing blockers or duplicate findings.
Routed skills: review-comet-pr and review-comet-expression-pr. Read AGENTS.md and the relevant contributor guidance.
Summary
- Prior state and problem: Registration accepted signatures that could not be serialized for native execution. Subsequent calls failed because Spark cannot evaluate native UDFs. The existing exception also attributed every Spark evaluation to operator fallback, overlooking evaluation during planning and range-bound sampling.
- Design approach: Validate argument and return types through
QueryPlanSerde.serializeDataTypebefore loading the library or replacing the session function. Explain the existing evaluation limitations in the exception, scaladoc, and Rust UDF guide. - Correctness: The guard checks both arguments and the return type, including recursively unsupported types inside arrays, maps, and structs. Rejection precedes registration side effects. Spark source confirms that
ConvertToLocalRelationevaluates projections, partition pruning evaluates predicates, and range partitioning evaluates sort-key projections. These paths explain the documented failures without changing evaluation semantics. - Compatibility analysis: Compared the relevant Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Their relevant evaluation behavior is consistent.
ExpectsInputTypescontinues to enforce the existing signature contract. No ABI, kernel, or shim changes occur. The latest release branch,branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, contains none of the touched UDF API files, so this intentionally adjusts an unreleased experimental feature. - Key design decisions: Reusing the serializer avoids a separate type whitelist. Validation happens before library loading and function replacement. The existing exception class remains intact while its message becomes more accurate.
- Implementation sketch: A signature scan finds the first unserializable type and raises
IllegalArgumentExceptionnaming it. Three tests cover planning-time evaluation overVALUES, sorting a projected UDF result, and refusal of anObjectTypeargument without installing a function. - Performance: The added work is a registration-time traversal and protobuf construction for the signature. It adds no per-row or per-batch processing. No evidence-backed performance regression was identified, and no benchmark improvement is claimed.
- Design: The change is narrow and easy to follow. It rejects an unusable signature earlier while documenting existing Spark execution constraints. It does not introduce optimizer or execution-routing changes.
- Abstraction & complexity: The existing serializer is an appropriate shared check for representability. No new abstraction is needed. This check does not replace the library's existing validation of whether a particular kernel supports a signature.
- Behavioral changes worth calling out: Unserializable signatures now fail during registration. Spark-side evaluation retains the same exception type with expanded diagnostic text. The guide's projected-column sorting workaround is exercised by a passing integration test.
- Suggested improvements: No further change meeting the P1/P2 reporting bar was identified.
Exact-head CI: run 37797850156 passed Required Checks, all four Linux Spark 4.1 test groups, native build/tests, and the Spark-profile compilation/lint matrix. The expressions-job log confirms all 71 CometNativeUdfSuite tests passed, including all three additions. Its overall result was 2,198 successful tests, zero failures, one canceled, and twelve ignored. Spark SQL suites, macOS, and other Spark runtime profiles were skipped. The separate label run skipped its substantive suites.
Validation limits: No local JVM/native suite or full Spark SQL shard was run. This checkout lacks built artifacts, and the environment encountered disk exhaustion during inspection. Runtime evidence comes from the exact-head CI logs, supplemented by source comparison across supported Spark versions. git diff --check passed, and project files remain unchanged.
Which issue does this PR close?
No issue. These are follow-ups to #4459, split out of #6697 so they can land without its vectorized JVM UDF API.
Rationale for this change
Two gaps in the native UDFs that #4459 added:
CometNativeUDF.registeraccepts a signature with a type Comet has no native representation for, such as a user-defined type. Every call then falls back to Spark, which cannot evaluate a native UDF, so the query fails at run time, far from the registration that caused it.VALUES, in a filter on partition columns, and when it samples a global sort's keys to choose range bounds. In those cases Comet runs every operator, yet the error says that Comet did not take the operator holding the call and points to the extended explain output, which shows no fallback. The scaladoc and the Rust UDF guide say the same.What changes are included in this PR?
CometNativeUDF.registerrefuses a signature with an argument or return type thatQueryPlanSerde.serializeDataTypecannot serialize, before it loads the library, and the error names the type.CometUdfNotEvaluatedException's message gives both causes: an operator Comet did not take, whose reason is in the extended explain output, or Spark evaluating the call while it plans the query, in the three cases above.CometNativeUDF.registerandNativeUdfCallsays the same.How are these changes tested?
Three new tests in
CometNativeUdfSuite:VALUESfails while Spark plans the query, with the error naming the planning-time casesObjectTypeargument is refused and installs no functionRemoving the type check fails the registration test, and taking the planning-time cases out of the error fails the
VALUEStest.CometNativeUdfSuitepasses on Spark 4.1 (71 tests).