Skip to content

Sliced booleans nested in structs, and sliced inputs to the JVM UDF bridge, reach the JVM misaligned and return wrong results #6288

Description

@andygrove

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:ffiArrow FFI / JNI boundarybugSomething isn't workingcorrectnesspriority:criticalData corruption, silent wrong results, security issues

Type

No type

Projects

No projects

    Milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions