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]

Reply via email to