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]