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 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 takes 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::slices.
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:
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:
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.
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.offsetat every level.ArrayImporter.doImportnever readssnapshot.offsetin 18.3.0 or 19.0.0, and apache/arrow-java#88 is still open. arrow-rs folds a slice into the buffers for most types, so a slicedInt64Array,StringArrayorStructArrayexports offset 0.BooleanArrayis the exception. It keeps its bit offset inArrayData::offset, andalign_nullskeeps 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_outputinnative/core/src/execution/jni_api.rstakes a column whenarray_ref.offset() != 0. Two paths are still exposed:prepare_outputsends it as is, and the JVM misreads the boolean children. The repros below go throughCometColumnarToRow. Broadcast serialization, the in-memory cache and JVM columnar shuffle read the same imported vectors, so they should see the same wrong values.JvmScalarUdfExpr::evaluateinnative/spark-expr/src/jvm_udf/mod.rsexports its argument arrays withFFI_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 booleanScalaUDFarguments and any codegen-dispatched expression that reads a boolean column, such asregexp_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 asliceof the emitted groups once a partition has more thanspark.comet.batchSizegroups), and any other operator that emitsRecordBatch::slices.array<boolean>andmap<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
mainat31b38196dwith Spark 4.1.3 and with Spark 3.5 / Scala 2.12:Over aggregate output.
spark.sql.shuffle.partitions=1only puts more than 8192 groups in one partition; with bigger data the default partitioning does the same:SELECT k, regexp_replace(IF(b, s, 'zz'), '1', 'y') FROM (SELECT k, b, s, count(*) FROM g GROUP BY k, b, s), withs = cast(k AS string), picks the wrong branch the same way.All of these plans are fully native, and
checkSparkAnswerfails on each. For the first query Comet returns[[true,17]],[[false,21]],[[true,24]]where Spark returns[[false,17]],[[true,21]],[[false,24]]. TheLIMIT ... OFFSETquery over anarray<boolean>ormap<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.offseton import.Additional context
A possible fix is to replace the top-level
array_ref.offset() != 0check with a recursive one that visitschild_dataand dictionary values. Any column where it finds a non-zero offset gets normalized beforemove_to_spark, andJvmScalarUdfExpr::evaluatedoes the same for its inputs. The existingcopy_array(aMutableArrayDatacopy) 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 booleanScalaUDFargument, each past the first batch.