You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
Support custom vectorized Java UDFs that operate on Arrow vectors #6694
#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.
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 feat: custom Rust UDFs via arrow-ffi [experimental] #4459.
Planning. In feat: custom Rust UDFs via arrow-ffi [experimental] #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 feat: custom Rust UDFs via arrow-ffi [experimental] #4459, argument types that differ from the registered ones are refused at planning time.
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 feat: custom Rust UDFs via arrow-ffi [experimental] #4459's SDK does. This overlaps Tighten CometUDF API with input/return type validation at registration #4173.
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:
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 feat: custom Rust UDFs via arrow-ffi [experimental] #4459 turned down, although Arrow Java's vector API changes far less often than DataFusion's.
Pass the C Data Interface, as feat: custom Rust UDFs via arrow-ffi [experimental] #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.
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.
sounds great, we prob need document well how writing row level UDFs is different from vectorized, additionally we can have the agent skill that can help to rewrite existing UDFs
Would be interesting to see limits of this approach
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 (
CometScalaUDFandCometScalaUDFCodegen). 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:
JvmScalarUdfinexpr.protonames 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.CometUdfBridgeimports them as Arrow Java vectors and callsCometUDF.evaluate(inputs: Array[ValueVector], numRows: Int). It keeps one instance per class per task attempt and installs the task'sTaskContextand context ClassLoader, so a class from--jarsresolves. It checks that the result is aFieldVectorwithnumRowsrows and exports it.The dispatcher is the only
CometUDFimplementation, 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
CometScalaUDF.convertchecks the registry before falling through to the dispatcher. A Java entry would emitJvmScalarUdfnaming 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. Inadd_one(abs(x)),abswould run natively and only its result would cross into the JVM. As in feat: custom Rust UDFs via arrow-ffi [experimental] #4459, argument types that differ from the registered ones are refused at planning time.registerAll(Native UDFs: lift the 4-argument cap and addregisterAll#6177).return_typetoday. The native side imports it with whatever schema the JVM exported, so a UDF that returns anIntVectorfor a declaredLongTypehands DataFusion anInt32column where it was promisedInt64. The dispatcher gets the type right by construction, but user code won't always. The bridge should check it and name both types, as feat: custom Rust UDFs via arrow-ffi [experimental] #4459's SDK does. This overlaps Tighten CometUDF API with input/return type validation at registration #4173.evaluatereceives 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 (Register CometArrowAllocator as a Spark MemoryConsumer for JVM-UDF dispatch #4174, feat: account JVM UDF Arrow allocations in Spark task memory #5027).CometUDF's scaladoc: literal arguments arrive as length-1 vectors, the result must havenumRowsrows, 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 toorg.apache.arrow.vector.ValueVector, butspark/pom.xmlrelocates Arrow in the published jar.javapon the shadedcomet-sparkjar shows the interface a user would actually compile against:A UDF written against stock Arrow Java doesn't implement it. The options:
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 thedatafusion-fficoupling feat: custom Rust UDFs via arrow-ffi [experimental] #4459 turned down, although Arrow Java's vector API changes far less often than DataFusion's.org.apache.arrow.cunrelocated, linked against the relocated vector classes. With Comet onextraClassPathand Spark's default parent-first class loading, a user jar'sorg.apache.arrow.c.Dataresolves to Comet's copy, and itsimportVectortakesorg.apache.comet.shaded.arrow.memory.BufferAllocator. That collision would have to be solved first.ColumnVector, which Comet'sCometVectoralready 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
CometNativeUdfSuiteis the template for tests, including the throwing stub that turns a silent fallback into a test failure.CometUDFregistration in May and was closed when the codegen dispatcher (feat(experimental): ScalaUDF and Java UDF support via Janino codegen #4267) landed. Its registration modes, including pairing with a row-based Spark UDF, are worth another look.