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]
