andygrove opened a new issue, #6694:
URL: https://github.com/apache/datafusion-comet/issues/6694
### What is the problem the feature request solves?
#4459 adds custom scalar UDFs written in Rust, which run inside the native
plan and operate on whole Arrow arrays. There is no equivalent for Java or
Scala.
A Java or Scala UDF does stay in the Comet pipeline today, through the
codegen dispatcher (`CometScalaUDF` and `CometScalaUDFCodegen`). But the
dispatcher still calls the user's function once per row, converting values to
and from the function's Scala/Java types on every call. A function that could
work on a column at a time, such as a tight loop over a primitive buffer or one
call per batch into a batch-oriented Java library, has no way to do so. The
only vectorized option is to rewrite it in Rust, which means building a native
library per platform and putting it on every executor (#6176). A Java class
would ship in the application jar like any other UDF.
### Describe the potential solution
#### Most of the execution path already exists
The dispatcher sits on a general JVM UDF bridge, and nothing in the bridge
is specific to it:
- `JvmScalarUdf` in `expr.proto` names a class and carries the argument
expressions, return type and nullability.
- `JvmScalarUdfExpr` (`native/spark-expr/src/jvm_udf/mod.rs`) evaluates the
arguments natively, exports them over the Arrow C Data Interface and calls into
the JVM.
- `CometUdfBridge` imports them as Arrow Java vectors and calls
`CometUDF.evaluate(inputs: Array[ValueVector], numRows: Int)`. It keeps one
instance per class per task attempt and installs the task's `TaskContext` and
context ClassLoader, so a class from `--jars` resolves. It checks that the
result is a `FieldVector` with `numRows` rows and exports it.
The dispatcher is the only `CometUDF` implementation, and it is the only
class the serde ever names. What's missing is the user-facing half.
#### Mirror #4459, with a class in place of a library
```scala
CometJvmUDF.register( // name to be decided
spark,
name = "add_one",
udfClass = classOf[AddOne], // implements the vectorized UDF interface
inputTypes = Seq(LongType),
returnType = LongType)
```
- **Registration.** Validate on the driver that the class loads, implements
the interface and has a public no-arg constructor. Then write the registry
entry, and install the catalog stub last, in the same order as #4459.
- **Planning.** In #4459, `CometScalaUDF.convert` checks the registry before
falling through to the dispatcher. A Java entry would emit `JvmScalarUdf`
naming the user's class, with each argument serialized as a native expression.
That differs from the dispatcher, which compiles the whole argument subtree
into its JVM kernel. In `add_one(abs(x))`, `abs` would run natively and only
its result would cross into the JVM. As in #4459, argument types that differ
from the registered ones are refused at planning time.
- **One registry for both paths.** Then the follow-ups already filed against
#4459 get solved once: session scoping (#5294), a registered UDF answering
calls to an ordinary Scala UDF with the same name (#5295), and the 4-argument
stub cap and `registerAll` (#6177).
- **Check the result type.** Nothing compares the result with the declared
`return_type` today. The native side imports it with whatever schema the JVM
exported, so a UDF that returns an `IntVector` for a declared `LongType` hands
DataFusion an `Int32` column where it was promised `Int64`. The dispatcher gets
the type right by construction, but user code won't always. The bridge should
check it and name both types, as #4459's SDK does. This overlaps #4173.
- **Pass an allocator.** `evaluate` receives no allocator, so a UDF borrows
the import allocator from its inputs, and a zero-argument UDF has nothing to
borrow from. Ideally the allocator it gets is charged to the Spark task (#4174,
#5027).
- **Document the contract in the user guide.** Today it lives only in
`CometUDF`'s scaladoc: literal arguments arrive as length-1 vectors, the result
must have `numRows` rows, and each task gets its own instance, called by one
thread at a time.
One choice #4459 didn't face: when Comet doesn't take the operator, a Rust
UDF can only fail, which is why its stub throws. A Java UDF could still run on
Spark. That could go through a row-based implementation registered alongside
it, as #4233 offered, or through an adapter over the vectorized one. The
alternative is to keep #4459's fail-loud stub, which also makes a silent
fallback fail the tests.
#### Open question: Comet relocates Arrow
`CometUDF`'s source refers to `org.apache.arrow.vector.ValueVector`, but
`spark/pom.xml` relocates Arrow in the published jar. `javap` on the shaded
`comet-spark` jar shows the interface a user would actually compile against:
```
public interface org.apache.comet.udf.CometUDF {
public abstract org.apache.comet.shaded.arrow.vector.ValueVector
evaluate(org.apache.comet.shaded.arrow.vector.ValueVector[], int);
}
```
A UDF written against stock Arrow Java doesn't implement it. The options:
1. **Compile against the relocated classes**
(`org.apache.comet.shaded.arrow.*`). The bridge works unchanged. The cost is
that user source is tied to Comet's Arrow Java version and relocation prefix.
That's the Java counterpart of the `datafusion-ffi` coupling #4459 turned down,
although Arrow Java's vector API changes far less often than DataFusion's.
2. **Pass the C Data Interface**, as #4459 does, and let the UDF import with
its own Arrow Java. This decouples the versions, but Comet's jar deliberately
leaves `org.apache.arrow.c` unrelocated, linked against the relocated vector
classes. With Comet on `extraClassPath` and Spark's default parent-first class
loading, a user jar's `org.apache.arrow.c.Data` resolves to Comet's copy, and
its `importVector` takes
`org.apache.comet.shaded.arrow.memory.BufferAllocator`. That collision would
have to be solved first.
3. **Keep Arrow out of the API.** Pass the UDF Spark's public
`ColumnVector`, which Comet's `CometVector` already implements over Arrow with
no copy. Add a Comet-owned writer for the result, since Spark has no public
writable vector. Nothing is tied to Comet's Arrow version, but the UDF reads
through accessors instead of Arrow buffers.
Option 1 is the smallest first step. It also fits the experimental footing
of #4459, where a UDF is rebuilt against each Comet release.
### Additional context
- #4459 is the model. Its `CometNativeUdfSuite` is the template for tests,
including the throwing stub that turns a silent fallback into a test failure.
- #4233 proposed user-facing `CometUDF` registration in May and was closed
when the codegen dispatcher (#4267) landed. Its registration modes, including
pairing with a row-based Spark UDF, are worth another look.
- User code is arbitrary, so these bite harder for user UDFs than for the
dispatcher: #4175 (task cancellation) and #6293 (a blocking JVM call holds a
Tokio worker for as long as it runs).
- #5596 (SPARK-55278): its point about reusing the SPIP's registration
vocabulary applies here as well.
--
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]