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]

Reply via email to