What is the problem the feature request solves?
#6715 sends Spark's null guard around a Scala UDF to the codegen dispatcher together with the UDF, so the null check is a branch in the kernel's loop rather than a DataFusion CASE that filters the batch for each branch and merges the results. It matches only the exact If that HandleNullInputsForUDF builds, if (isnull(a) or ...) null else f(knownnotnull(a), ...), so these three shapes still run the guard natively. In each one, with plusOne = (x: Long) => x + 1 over a nullable BIGINT column a, extended explain on Spark 4.1 lists if and isnull as native:
- Nested primitive UDFs.
plusOne(plusOne(a)) optimizes to IF(a IS NULL, NULL, plusOne(knownnotnull(if (isnull(a)) null else plusOne(knownnotnull(a))))). The inner guard is part of the outer UDF's argument and already runs in its kernel. The outer guard's predicate starts as isnull(<inner guard>), and PushFoldableIntoBranches pushes the isnull into the inner guard's branches, which simplifies to a IS NULL. That predicate no longer matches the UDF's argument.
- A foldable operation on the result.
PushFoldableIntoBranches also pushes the + 1 in plusOne(a) + 1 into the branches, giving IF(a IS NULL, NULL, plusOne(knownnotnull(a)) + 1). The false branch is an Add rather than the UDF, and add runs natively too.
- Guards that users write. Spark adds no guard for a UDF with boxed parameters, but queries often write one, such as
IF(a IS NULL, NULL, plusOneBoxed(a)). The argument has no KnownNotNull, so it does not match.
Describe the potential solution
The kernel runs Spark's own generated code for every expression in it, so taking more guard shapes does not change results. The exact match only keeps work that DataFusion vectorizes out of the row loop. Each shape needs a narrow extension of the match in CometScalaUDF.isNullGuard:
- Accept a predicate whose
IsNull checks are on expressions inside the guarded arguments, such as a under the inner guard.
- Accept a false branch that applies operations with foldable operands to the UDF, which is what
PushFoldableIntoBranches produces.
- Accept
IsNull checks on any argument of the UDF, not only the KnownNotNull ones.
Each extension has to keep two things #6715 does: nondeterministic arguments stay native, and the guard stays native when a node it would pull into the kernel has spark.comet.expression.<name>.enabled=false.
Additional context
@mbutrovich asked for this issue on #6715 (#6715 (review)). For a sense of what the native guard costs, #6715 measured SELECT max(f(c)) over 4M rows at batch size 8192: dispatching the guard with the UDF cut the time above the max(c) floor from 42 ms to 26.5 ms.
What is the problem the feature request solves?
#6715 sends Spark's null guard around a Scala UDF to the codegen dispatcher together with the UDF, so the null check is a branch in the kernel's loop rather than a DataFusion
CASEthat filters the batch for each branch and merges the results. It matches only the exactIfthatHandleNullInputsForUDFbuilds,if (isnull(a) or ...) null else f(knownnotnull(a), ...), so these three shapes still run the guard natively. In each one, withplusOne = (x: Long) => x + 1over a nullableBIGINTcolumna, extended explain on Spark 4.1 listsifandisnullas native:plusOne(plusOne(a))optimizes toIF(a IS NULL, NULL, plusOne(knownnotnull(if (isnull(a)) null else plusOne(knownnotnull(a))))). The inner guard is part of the outer UDF's argument and already runs in its kernel. The outer guard's predicate starts asisnull(<inner guard>), andPushFoldableIntoBranchespushes theisnullinto the inner guard's branches, which simplifies toa IS NULL. That predicate no longer matches the UDF's argument.PushFoldableIntoBranchesalso pushes the+ 1inplusOne(a) + 1into the branches, givingIF(a IS NULL, NULL, plusOne(knownnotnull(a)) + 1). The false branch is anAddrather than the UDF, andaddruns natively too.IF(a IS NULL, NULL, plusOneBoxed(a)). The argument has noKnownNotNull, so it does not match.Describe the potential solution
The kernel runs Spark's own generated code for every expression in it, so taking more guard shapes does not change results. The exact match only keeps work that DataFusion vectorizes out of the row loop. Each shape needs a narrow extension of the match in
CometScalaUDF.isNullGuard:IsNullchecks are on expressions inside the guarded arguments, such asaunder the inner guard.PushFoldableIntoBranchesproduces.IsNullchecks on any argument of the UDF, not only theKnownNotNullones.Each extension has to keep two things #6715 does: nondeterministic arguments stay native, and the guard stays native when a node it would pull into the kernel has
spark.comet.expression.<name>.enabled=false.Additional context
@mbutrovich asked for this issue on #6715 (#6715 (review)). For a sense of what the native guard costs, #6715 measured
SELECT max(f(c))over 4M rows at batch size 8192: dispatching the guard with the UDF cut the time above themax(c)floor from 42 ms to 26.5 ms.