Comet executes Spark's Scala and Java scalar user-defined functions (UDFs) within the Comet pipeline. The presence of a UDF does not force the enclosing operator out of the Comet pipeline; surrounding Rust-implemented operators stay in the pipeline.
This page covers Spark's ScalaUDF (Scala udf(...), spark.udf.register(...) over Scala or Java functional interfaces, and SQL CREATE FUNCTION ... AS 'com.example.MyUDF'). Other UDF kinds (Python / Pandas, Hive, aggregate) are out of scope and continue to fall back to Spark.
This feature is enabled by default. Set spark.comet.exec.scalaUDF.codegen.enabled to false to route plans containing a ScalaUDF back to Spark for the enclosing operator.
| Key | Default | Description |
|---|---|---|
spark.comet.exec.scalaUDF.codegen.enabled | true | When true, eligible ScalaUDFs run in the Comet pipeline. When false, the enclosing operator falls back to Spark. |
udf(...), spark.udf.register(...) (Scala or Java functional interfaces), or SQL CREATE FUNCTION ... AS 'com.example.MyUDF'.Boolean, Byte, Short, Int, Long, Float, Double, Decimal, String, Binary, Date, Timestamp, TimestampNTZ.ArrayType, StructType, MapType.myUdf(upper(s)) runs as one unit in the Comet pipeline).transform, filter, exists, aggregate, zip_with, map_filter, map_zip_with, etc.) inside the argument tree.ScalaAggregator, TypedImperativeAggregate, the legacy UserDefinedAggregateFunction).@udf and Pandas @pandas_udf.GenericUDF and SimpleUDF.CalendarIntervalType, NullType, and UserDefinedType arguments and return types. UDT-typed columns fall back to Spark; to keep execution in the Comet pipeline, store and read the underlying representation directly (e.g. write MLlib Vector outputs as Struct<type: Byte, size: Int, indices: Array<Int>, values: Array<Double>> rather than VectorUDT).spark.sql.codegen.maxFields (default 100). Comet refuses these at plan time and the operator falls back to Spark.When a UDF is rejected, the reason surfaces through Comet's standard fallback diagnostics; the query still runs on Spark.
rand, uuid, monotonically_increasing_id) produce per-partition sequences consistent with Spark.TaskContext.get() inside the user function returns the driving Spark task's context.CodeGenerator cache, so structurally identical queries across a session share the compiled class.