andygrove opened a new pull request, #6697:
URL: https://github.com/apache/datafusion-comet/pull/6697

   ## Which issue does this PR close?
   
   Closes #6694.
   
   ## Rationale for this change
   
   #4459 adds custom scalar UDFs written in Rust. There is no vectorized 
equivalent for Java or Scala: an ordinary Java or Scala UDF stays in the Comet 
pipeline through the codegen dispatcher, but the dispatcher calls it once per 
row and converts every value to and from the function's types. A function that 
could work on a whole column at once has no way to, short of being rewritten in 
Rust and shipped as a native library for every platform.
   
   Most of the execution path already exists. `JvmScalarUdf` names a class and 
carries argument expressions, `JvmScalarUdfExpr` evaluates the arguments 
natively and exports them over the Arrow C Data Interface, and `CometUdfBridge` 
calls `CometUDF.evaluate` with one instance per class per task. The dispatcher 
has been the only `CometUDF` implementation and the only class the serde names. 
This PR adds the user-facing half, mirroring #4459 with a class in place of a 
library.
   
   ## What changes are included in this PR?
   
   - **`CometJvmUDF.register(spark, name, udfClass, inputTypes, returnType, 
deterministic = true)`**, plus an overload taking a `java.util.List` for Java 
callers. It is experimental and deliberately not annotated `@Public`. It first 
checks on the driver that the class is concrete and public and has a public 
no-argument constructor, which is what the bridge needs to instantiate it. Then 
it writes the registry entry and installs the catalog stub last, in the same 
order as #4459.
   - **`CometUdfRegistry`**: a driver-side registry keyed by function name. Its 
entries are a sealed `UdfMetadata` with a single case today, `JvmUdfMetadata`. 
#4459's native entries can join as a second case rather than living in a second 
registry (see Follow-on work).
   - **Catalog stub**: built from Spark's Java `UDF0`..`UDF4`, so Spark inserts 
no casts for its arguments. It throws `CometUdfNotEvaluatedException` if Spark 
ever evaluates it, so a silent fallback becomes a visible failure, as in #4459. 
The arity cap is #4459's four arguments (#6177).
   - **Planning**: `CometScalaUDF.convert` checks the registry before the 
dispatcher.
     - A registered name emits a `JvmScalarUdf` that names the user's class, 
with each argument serialized as its own native expression. In 
`add_one(abs(x))`, `abs` runs natively and only its result crosses into the 
JVM. The dispatcher, by contrast, compiles the whole argument tree into its JVM 
kernel.
     - A call whose argument types differ from the registered ones (nullability 
aside) is refused at planning time, naming both signatures. This is #4459's 
rule.
     - No proto change: `JvmScalarUdf` already carries everything.
   - **Result type check** in the native `JvmScalarUdfExpr`. A result used to 
be imported with whatever schema the JVM exported, so a UDF declared as 
`LongType` that returned an `IntVector` handed DataFusion an `Int32` column 
under an `Int64` schema. Each result is now compared with the declared return 
type:
     - **Relabelled**: differences that don't change the data. These are list 
and map child field names (Arrow Java's `ListVector` names its element 
`$data$`), nested nullability and metadata, and `Null`-typed children (a 
`ListVector` that only received empty lists has a `Null` element).
     - **Rejected**: anything else fails the query and names both types.
     - The check covers the dispatcher's results too. They already match by 
construction, and `CometCodegenSuite` passes unchanged.
   - **Arrow relocation**: option 1 from the issue. A UDF compiles against 
`org.apache.comet.shaded.arrow.*` from the published jar, which ties its source 
to Comet's relocation prefix and Arrow Java version. That fits the experimental 
footing; the user guide says so and says to rebuild per release.
   - **User guide**: a new page, `Vectorized Java and Scala UDFs`, linked from 
the Scala/Java UDF page. It covers:
     - how a vectorized UDF differs from a row-level one
     - the contract, which until now lived only in `CometUDF`'s scaladoc: 
literal arguments arrive as one-value vectors, nulls are the UDF's to handle, 
the result's row count and Arrow type, the allocator, and instance lifetime and 
threading
     - registration, the relocation, and limitations
   - **CI**: `CometJvmUdfSuite` is listed in both the Linux and macOS workflows.
   
   ### Follow-on work
   
   - **Allocator.** `evaluate` still receives no allocator, so a UDF borrows 
its arguments' allocator and a zero-argument UDF has none. #5027 changes 
`CometUDF.evaluate` to take an allocator charged to the task and adds 
`close()`. Whichever of the two PRs lands second updates the other's UDFs and 
docs. #5027 also corrects the threading contract, and the user guide here 
already follows the corrected version: `evaluate` can be called concurrently on 
one instance.
   - **One registry with #4459.** Whichever of #4459 and this PR lands second 
moves the other's entries into `CometUdfRegistry`. Session scoping (#5294), 
name collisions with ordinary Scala UDFs (#5295), and the arity cap and 
`registerAll` (#6177) can then be solved once for both.
   - **Running on Spark when Comet does not take the operator.** The stub fails 
loudly, as in #4459. An adapter that runs the vectorized UDF over one-row 
vectors could follow if a silent fallback turns out to be wanted.
   - User code is arbitrary, so #4175 (task cancellation) and #6293 (a blocking 
call holds a native execution thread) apply here with more force than to the 
dispatcher.
   
   ## How are these changes tested?
   
   `CometJvmUdfSuite` (39 tests). Like #4459's suite, it is self-guarding: the 
catalog stub throws if Spark evaluates the UDF, so a fallback fails a test 
instead of passing it. It covers:
   
   - argument order, nulls, and literal arguments in each position
   - arguments evaluated natively with the codegen dispatcher turned off
   - the UDF in a filter, a join condition, a grouping key, and a window 
`PARTITION BY`, each of which reaches the serde through a different operator
   - a round trip of 18 types, primitive and nested, each with a null row
   - lists built with Arrow Java's own writer, for both `containsNull` values; 
a batch of only empty lists; and a list combined with Comet's own arrays in 
both branch orders of `if`
   - a result of the wrong type (the error names both types and the class) and 
a result with the wrong row count
   - argument types other than the registered ones being refused
   - the stub failing with Comet disabled, and a nondeterministic registration
   - registration refusing an abstract class, a class without a no-argument 
constructor, an inner class, and more than four arguments
   - a Java UDF compiled at test time and visible only on 
`spark.executor.extraClassPath`, registered through the Java overload and 
called from a native parquet scan and from a range
   
   Native unit tests in `jvm_udf::result_type` (11) cover the relabelling and 
rejection rules. `CometCodegenSuite`, `CometCodegenHOFSuite`, 
`CometUdfBridgeSuite` and `CometScalaUDFClassLoaderSuite` pass unchanged with 
the result type check in place.
   
   `CometJvmUdfBenchmark` compares one function in three forms: an ordinary 
Scala UDF on Spark, the same UDF through the codegen dispatcher, and a 
vectorized UDF that loops over the Arrow buffers.
   
   Each query is `SELECT max(f(c)) FROM t` over 4M rows of a parquet `bigint` 
column, a tenth of them null. The last row runs the query without the UDF, so a 
UDF row's distance above it is what that form of the UDF costs. Best of runs, 
Apple M3 Ultra, JDK 17, Spark 4.1, release build:
   
   | Function | Spark | Comet, codegen dispatch | Comet, vectorized UDF | 
Comet, no UDF |
   | --- | --- | --- | --- | --- |
   | `x + 1` | 79 ms | 91 ms | 58 ms | 52 ms |
   | 64-bit mix (about 10 arithmetic ops) | 78 ms | 89 ms | 55 ms | 46 ms |
   
   Above the no-UDF floor, the vectorized UDF costs about 1.5 to 2 ns per row 
and the dispatched UDF about 10 ns per row. Neither form slows down from `x + 
1` to the mix, so at this size the cost is in calling the function rather than 
in its arithmetic. That per-row call is what a vectorized UDF saves. A function 
whose own work dominates gains correspondingly less, which bounds what 
rewriting a UDF in vectorized form can buy.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to