andygrove opened a new pull request, #6542:
URL: https://github.com/apache/datafusion-comet/pull/6542
## Which issue does this PR close?
Closes #6425.
This replaces #6455, which tried to fix the same issue inside the dispatcher.
## Rationale for this change
Since #5692, `CometStaticInvoke` and `CometInvoke` send every call they
don't otherwise handle to the codegen dispatcher. That includes calls to
DataSource V2 catalog functions, which are user code. Such a function can
return a `Decimal` at a different scale from the type it declares, or one that
doesn't fit it. Spark only corrects that when it writes a row, and any
expression around the call reads the value as returned. The dispatcher has to
write an Arrow vector of the declared type, so whatever reads its output sees a
different value.
#6455 rescaled the value in the dispatcher and tried to keep every consumer
of the call in the same kernel or in Spark. But where Spark writes rows depends
on whole-stage codegen, and each review round found another path: transitive
parents, projected aliases, `explode`, a dispatched `map(...)` around the call,
and `ApplyFunctionExpression`. In 1.0.0 these calls fell back to Spark, so this
PR goes back to that. The dispatcher's handling of `StaticInvoke` and `Invoke`
becomes an allow list instead of a catch-all.
## What changes are included in this PR?
The new `CometInvokeTargets` decides which calls the dispatcher may run. A
`StaticInvoke` or `Invoke` is allowed when the class it calls is part of Spark
(`org.apache.spark.sql.catalyst.*` or `org.apache.spark.unsafe.*`) and isn't a
`ScalarFunction`. An `Invoke` on a Catalyst value, such as Spark 4's
`is_valid_utf8` calling `isValid` on a `UTF8String`, is allowed too. An
`ApplyFunctionExpression` never is. The one exception is the predicate of a
typed `Dataset.filter`. Spark calls it through `scala.Function1` or
`FilterFunction`, and it returns a boolean, so there's nothing for a row write
to correct. #6497 covers that path, so it stays on the dispatcher.
`emitJvmCodegenDispatch` applies the check to the whole tree rather than
just its root, because the kernel runs every node under the dispatched
expression. That's what catches a DSv2 call nested inside `map(...)` or another
dispatched expression.
Spark's own lowerings keep dispatching, for example binary `lpad`/`rpad`,
`encode`/`to_binary`, `to_time` and the Spark 4 evaluator `Invoke`s. I checked
the `StaticInvoke` and `Invoke` targets in Spark 3.4.3 through 4.2.0, and every
function lowering calls into one of those two packages. Iceberg functions with
a native handler are unaffected. An Iceberg function with no handler, or nested
inside a dispatched expression, now falls back. The dispatcher's decimal writer
is unchanged. The user guide now says that other DSv2 functions run in Spark.
## How are these changes tested?
`CometCodegenSuite` registers a DSv2 function catalog whose functions return
scale-0 decimals for `DECIMAL(10, 2)`, one for each lowering (`Invoke`,
`StaticInvoke` and `ApplyFunctionExpression`). It runs the shapes from the
#6455 reviews:
- the plain result, `IS NULL` and a cast to string
- `abs(...)`, and array and struct access on the result
- `count`, `max` and `sum`
- a projected alias, `explode`, and `map(...)` around the call
- `DISTRIBUTE BY` on the call
Each one falls back with the new reason and matches Spark. A unit test pins
which calls the allow list accepts and declines, both at the root and nested
under `map(...)`. `CometIcebergSystemFunctionSuite` adds @sunchao's
`map_values(map('k', truncate(10, d)))[0]` case.
The #5573 and #5575 tests used test classes as `Invoke` targets, which the
dispatcher now declines. The #5573 test now uses a UDF that captures a
non-serializable object, and the #5575 test uses an `Invoke` on a `UTF8String`.
The unlisted `StaticInvoke` test in the Iceberg suite now expects a fallback.
With the check disabled, both new SQL tests fail. The DSv2 one returns
#6425's wrong answers (`0.03` for `3.00`, `1000000.00` where Spark returns
null).
Local runs:
- Spark 4.1: `CometCodegenSuite` and `CometIcebergSystemFunctionSuite` (121
tests), `CometExpressionSuite` (174), and `CometSqlFileTestSuite`,
`CometStringExpressionSuite`, `CometCodegenSourceSuite`,
`CometCodegenFuzzSuite` and the two Iceberg function extension and pushdown
suites (730).
- Spark 3.4, 3.5, 4.0 and 4.2: `CometCodegenSuite`,
`CometIcebergSystemFunctionSuite` and `CometStringExpressionSuite` (158, 158,
160 and 146). The cancellations are version guards, and there's no Iceberg
build for Spark 4.2.
--
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]