andygrove commented on code in PR #6455:
URL: https://github.com/apache/datafusion-comet/pull/6455#discussion_r4164524984


##########
spark/src/main/scala/org/apache/comet/codegen/CometBatchKernelCodegenOutput.scala:
##########
@@ -231,15 +231,39 @@ private[codegen] object CometBatchKernelCodegenOutput 
extends CometTypeShim {
       val set = if (nested) "setSafe" else "set"
       OutputEmit("", s"$targetVec.$set($idx, $source);")
     case dt: DecimalType =>
+      // Rescale to the declared type, and write null when the value does not 
fit, as Spark's
+      // `UnsafeRowWriter` and `UnsafeArrayWriter` do in `write(ordinal, 
Decimal, precision,
+      // scale)`. A Spark expression already produces its declared precision 
and scale, but a
+      // DSv2 function called through `Invoke` / `StaticInvoke` can return a 
`Decimal` of any
+      // scale (#6425). Like Spark's writers, this rescales the value in 
place, and leaves it
+      // untouched when it does not fit. Unlike them, it does not test 
`source` for null: the
+      // callers write null values themselves, and skip that test only for a 
type that is not
+      // nullable.
+      //
+      // The precision and scale test repeats `changePrecision`'s own fast 
path. It keeps the call
+      // off the common path, so the JIT can still scalar-replace the 
`Decimal` that an input
+      // getter allocates. With the bare call, passing a `DECIMAL(18, 2)` 
column through took
+      // about half as long again per row.
+      //
       // DecimalOutputShortFastPath: precision <= 18 fits in a signed long, so 
pass the unscaled
       // value to `setSafe(int, long)` and skip the BigDecimal allocation.
+      val dec = ctx.freshName("dec")
+      val (precision, scale) = (dt.precision, dt.scale)
       val write =
-        if (dt.precision <= Decimal.MAX_LONG_DIGITS) {
-          s"$targetVec.setSafe($idx, $source.toUnscaledLong());"
+        if (precision <= Decimal.MAX_LONG_DIGITS) {
+          s"$targetVec.setSafe($idx, $dec.toUnscaledLong());"
         } else {
-          s"$targetVec.setSafe($idx, $source.toJavaBigDecimal());"
+          s"$targetVec.setSafe($idx, $dec.toJavaBigDecimal());"
         }
-      OutputEmit("", write)
+      OutputEmit(
+        "",
+        s"""org.apache.spark.sql.types.Decimal $dec = $source;
+           |if (($dec.precision() == $precision && $dec.scale() == $scale) ||
+           |    $dec.changePrecision($precision, $scale)) {
+           |  $write
+           |} else {
+           |  $targetVec.setNull($idx);

Review Comment:
   @sunchao The retained-projection case is fixed in 
c655d630696579c8d7ec3eac784afdc4a7fd1261, with the branch-1.1 cherry-pick in 
76a0e1c25d2f50dc7cf21e9f49bb7192acf502e0 (#6490). I reproduced both reported 
queries on the previous heads: the fused Spark projections preserve non-null 
overflowing values until the outer consumer runs, whereas the inner Comet 
projection had already normalized them. Both ANSI modes showed the 
nullness/count mismatch.
   
   The planner now follows decimal-bearing DSv2 projected attributes and keeps 
the affected producer, consumer, and intervening operators in Spark. Falling 
back only at the outer consumer would be too late. This is deliberately 
conservative: provenance can continue through downstream decimal outputs even 
beyond a real row boundary, so those affected chains may lose some native 
execution. Native input scans, unrelated branches, and calls consumed within 
one dispatched expression tree remain eligible for Comet.
   
   The regression covers retained instance/static aliases, intermediate 
projections, filters, array/struct access, scale-sensitive string casts, and 
count/max, with ANSI and AQE on/off. It also checks whole-stage codegen 
disabled, where Spark materializes each projection. Existing direct-output and 
parquet write/read boundary tests still pass. Additional full-SQL probes 
matched final Spark results and AQE plans, including count-only consumers and 
exchange/UNION boundaries.
   
   Local dispatcher, generated-source, and planner suites passed: source Spark 
4.1 257 tests; branch-1.1 Spark 4.1 253 tests; branch-1.1 Spark 3.5 251 tests, 
with two expected version-guard cancellations. Fresh configured CI, including 
the full Spark 4.1 SQL matrix, is now being monitored; its verdict is still 
pending.
   
   CI follow-up: the strict Spark 3.5 compiler caught an implicitly discarded 
return value in the traversal. The one-line explicit discard is now in 
e1972d144bc0ed18d121569dfe0d58716b5cd235 and release 
facb668c8af7e9eabc26d0118f5cf8a1fb11ce62. Local strict Spark 3.5 test 
compilation passes, and the source257/release253 Spark 4.1 tests passed again. 
This does not change the fallback behavior; replacement CI runs are pending.



-- 
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