andygrove opened a new issue, #6288:
URL: https://github.com/apache/datafusion-comet/issues/6288

   ### Describe the bug
   
   Comet returns wrong values, and puts nulls on the wrong rows, when a sliced 
boolean array crosses from native to the JVM anywhere below the top level of a 
column, or as an input to the JVM UDF bridge.
   
   The cause is that Arrow Java's C Data import ignores `ArrowArray.offset` at 
every level. `ArrayImporter.doImport` never reads `snapshot.offset` in 18.3.0 
or 19.0.0, and 
[apache/arrow-java#88](https://github.com/apache/arrow-java/issues/88) is still 
open. arrow-rs folds a slice into the buffers for most types, so a sliced 
`Int64Array`, `StringArray` or `StructArray` exports offset 0. `BooleanArray` 
is the exception. It keeps its bit offset in `ArrayData::offset`, and 
`align_nulls` keeps the validity bitmap at the same offset, so the JVM reads 
both the values and the nulls from bit 0.
   
   #2051 handled this for top-level columns. `prepare_output` in 
`native/core/src/execution/jni_api.rs` `take`s a column when 
`array_ref.offset() != 0`. Two paths are still exposed:
   
   - **A boolean nested in a struct.** A struct column has offset 0 even when 
its children are sliced, so `prepare_output` sends it as is, and the JVM 
misreads the boolean children. The repros below go through 
`CometColumnarToRow`. Broadcast serialization, the in-memory cache and JVM 
columnar shuffle read the same imported vectors, so they should see the same 
wrong values.
   - **Inputs to the JVM UDF bridge.** `JvmScalarUdfExpr::evaluate` in 
`native/spark-expr/src/jvm_udf/mod.rs` exports its argument arrays with 
`FFI_ArrowArray::new(&arr.to_data())` and no normalization at all. So a 
top-level sliced boolean, or a struct argument with sliced boolean children, 
reaches the UDF misaligned. That covers boolean `ScalaUDF` arguments and any 
codegen-dispatched expression that reads a boolean column, such as 
`regexp_replace(IF(b, s, 'zz'), '1', 'y')`. The dispatcher is on by default.
   
   Native slices are common. Examples are `ORDER BY ... LIMIT n OFFSET m`, 
DataFusion's grouped hash aggregate (every output batch after the first is a 
`slice` of the emitted groups once a partition has more than 
`spark.comet.batchSize` groups), and any other operator that emits 
`RecordBatch::slice`s.
   
   `array<boolean>` and `map<string, boolean>` are not affected, because 
slicing a list or a map moves its offsets buffer and leaves the values child 
whole.
   
   ### Steps to reproduce
   
   Default configs, reproduced on `main` at `31b38196d` with Spark 4.1.3 and 
with Spark 3.5 / Scala 2.12:
   
   ```scala
   spark.range(0, 1000)
     .selectExpr(
       "id AS v",
       "named_struct('x', id % 7 = 0, 'y', id) AS st",
       "named_struct('x', IF(id % 3 = 0, NULL, id % 7 = 0)) AS stn")
     .write.parquet(path)
   spark.read.parquet(path).createOrReplaceTempView("t")
   
   // CometTakeOrderedAndProjectExec(limit=57, offset=17): st.x is wrong on 
most rows.
   sql("SELECT st FROM t ORDER BY v LIMIT 40 OFFSET 17").collect()
   // The nulls land on the wrong rows too.
   sql("SELECT stn FROM t ORDER BY v LIMIT 40 OFFSET 17").collect()
   ```
   
   Over aggregate output. `spark.sql.shuffle.partitions=1` only puts more than 
8192 groups in one partition; with bigger data the default partitioning does 
the same:
   
   ```scala
   spark.range(0, 40000)
     .selectExpr("id % 20000 AS k", "(id % 20000) % 3 = 0 AS b")
     .write.parquet(path2)
   spark.read.parquet(path2).createOrReplaceTempView("g")
   spark.udf.register("flip", (x: Boolean) => !x)
   
   // Nested boolean from prepare_output: wrong after the first 8192 rows.
   sql("SELECT named_struct('b', b, 'k', k) FROM (SELECT k, b, count(*) FROM g 
GROUP BY k, b)").collect()
   // JVM UDF bridge input: flip(b) is computed from the wrong rows.
   sql("SELECT k, flip(b) FROM (SELECT k, b, count(*) FROM g GROUP BY k, 
b)").collect()
   ```
   
   `SELECT k, regexp_replace(IF(b, s, 'zz'), '1', 'y') FROM (SELECT k, b, s, 
count(*) FROM g GROUP BY k, b, s)`, with `s = cast(k AS string)`, picks the 
wrong branch the same way.
   
   All of these plans are fully native, and `checkSparkAnswer` fails on each. 
For the first query Comet returns `[[true,17]]`, `[[false,21]]`, `[[true,24]]` 
where Spark returns `[[false,17]]`, `[[true,21]]`, `[[false,24]]`. The `LIMIT 
... OFFSET` query over an `array<boolean>` or `map<string, boolean>` column 
matches Spark.
   
   ### Expected behavior
   
   Results match Spark. Native should hand the JVM arrays whose offsets are 
zero at every level, until Arrow Java honors `ArrowArray.offset` on import.
   
   ### Additional context
   
   A possible fix is to replace the top-level `array_ref.offset() != 0` check 
with a recursive one that visits `child_data` and dictionary values. Any column 
where it finds a non-zero offset gets normalized before `move_to_spark`, and 
`JvmScalarUdfExpr::evaluate` does the same for its inputs. The existing 
`copy_array` (a `MutableArrayData` copy) produces offset 0 at every level. A 
cheaper version would re-slice only the boolean value and validity buffers and 
leave the rest zero-copy. Since the aggregate case hits every batch after the 
first, the cheaper version is probably worth it.
   
   Tests should cover a sliced boolean under a struct (nullable and not), 
`list<struct<boolean>>` and a boolean `ScalaUDF` argument, each past the first 
batch.
   


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