Conversation
…he JVM Arrow Java's C Data import ignores ArrowArray.offset at every level (apache/arrow-java#88). arrow-rs folds a slice into the buffers for every Spark type except boolean, which keeps its bit offset, so a sliced boolean reached the JVM read from bit 0. The fix for apache#2051 only checked the top level of each output column, so a boolean inside a sliced struct, list or map came back wrong, and inputs to the JVM UDF bridge were never checked at all. Add zero_offsets to datafusion-comet-common. It returns the array unchanged when every offset is already zero, and otherwise re-slices boolean bitmaps to start at bit 0 at every level, sharing every other buffer. move_to_spark now applies it to every executed batch and decoded shuffle block, replacing the top-level take in prepare_output, and JvmScalarUdfExpr applies it to the bridge inputs. Closes apache#6288.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Top-level normalization missed sliced boolean children and JVM UDF inputs, allowing values and nulls to reach the JVM at incorrect bit positions.
- Design approach: Centralize recursive offset normalization in
zero_offsetsand apply it at both native-to-JVM export paths. - Correctness / compatibility analysis: Checked Arrow Java 18.3.0 import behavior, Arrow Rust 59.3.0 bitmap handling, and relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No introduced P1/P2 issues found within this review.
- Key design decisions: The common path borrows unchanged data. Boolean normalization replaces the index array and
takeoperation with bitmap slicing. Existing FFI ownership remains intact, and normal non-boolean buffers remain shared. No quantified performance claim was validated. - Implementation sketch: The shared helper handles children and dictionary values.
move_to_sparkandJvmScalarUdfExprinvoke it, with Rust and Scala regression coverage. - Behavioral changes worth calling out: Sliced nested booleans retain their values and null placement through limits, aggregates, collected structs and UDF inputs.
- Suggested improvements: None meeting the P1/P2 reporting threshold.
Reviewed full SHA 0c9f24941027785a9ea3a957c247d89f5e6b4f9b against e897f8ab45dec8fc27c79339c65eeaa7bd60e198, including all 14 differing files. The six additional memory/logging files reflect base-only commit #6250 being absent from this head. Their differences were inspected. The PR remains non-draft, and the snapshot and live discussion contain no reviews, comments or threads.
Routed skills: review-comet-pr, review-comet-ffi-pr, review-comet-memory-pr, and review-comet-expression-pr.
Exact-head CI: 20 checks succeeded, including Rust tests, native/JVM builds, scans, shuffle and TPC-H. Execution, expression and TPC-DS jobs remain running. Thirteen checks were skipped, including Spark SQL, Iceberg and macOS coverage. No failures were reported. Run: https://github.com/apache/datafusion-comet/actions/runs/36460041038
Validation: All seven local ffi_offsets tests passed. A disposable probe passed 5,460 independent value/validity offset combinations. Local core export tests were blocked by missing jni.h; the HDFS-disabled retry reached its 180-second compilation limit before executing tests. JVM/Spark suites and benchmarks were not run locally. The working tree remains unchanged.
Which issue does this PR close?
Closes #6288.
Rationale for this change
Arrow Java's C Data import ignores
ArrowArray.offsetat every level (apache/arrow-java#88). arrow-rs folds a slice into the buffers for every Spark type except boolean, which keeps its bit offset, so the JVM reads a sliced boolean's values and nulls from bit 0. The fix for #2051 only checked the top level of each output column. A boolean inside a sliced struct, list or map still came back wrong, and the inputs to the JVM UDF bridge were never checked at all. The issue has repros throughLIMIT ... OFFSET, a grouped aggregate emitting more than one batch, a booleanScalaUDFand a dispatchedregexp_replace.There is one more path that the issue doesn't list.
collect_listreturns a single run of rows as a slice, because arrow'sconcatof one array returns that array's slice, and that slice keeps its offsets. So a group whose only run starts past row 0 produced alist<struct<boolean>>with a sliced boolean two levels down.What changes are included in this PR?
zero_offsetsinnative/common/src/ffi_offsets.rsgives every level of an array offset 0, including struct children, list and map values, and dictionary values. When no level has an offset, which is almost always, it returns the array unchanged. A sliced boolean has its values bitmap re-sliced to start at bit 0. That shares the buffer when the offset is a whole number of bytes, and copies only the bitmap otherwise. Every other buffer is shared.FFI_ArrowArray::newalready realigns the validity bitmap to the exported offset.move_to_sparkapplies it, so it now covers every executed batch and decoded shuffle block. This replaces the top-leveltakeinprepare_output, and a top-level boolean now costs a bitmap slice instead of atake.JvmScalarUdfExpr::evaluateapplies it to the arrays it hands the JVM UDF bridge.ffi.mdgains a short "Array Offsets" section, so that a new native to JVM export path calls it too.How are these changes tested?
CometExecSuitecovers a struct sliced byLIMIT ... OFFSET(with and without nulls), anamed_structover sliced aggregate output, andcollect_listof structs.CometCodegenSuitecovers a booleanScalaUDFargument,regexp_replace(IF(b, s, 'zz'), ...)through the dispatcher, and a struct UDF input with a sliced boolean child. All six fail onmainwith wrong results and pass with this change, on Spark 4.1, 3.5 (Scala 2.12) and 3.4.zero_offsetscover a top-level boolean at byte-aligned and unaligned offsets, a struct child,list<struct<boolean>>, map values, dictionary values, the copy fallback for any other type with an offset, and the unchanged fast path. The six that rewrite something export throughFFI_ArrowArray, check that every level has offset 0, and import the array back. Amove_to_sparktest checks the export path itself. Removing the rewrite fails all six, keeping only a top-level check fails exactly the four nested ones, and zeroing the offset without re-slicing the bitmap fails every boolean case.CometExecSuite,CometCodegenSuite,CometUdfBridgeSuite,CometAggregateSuite,CometTopKSuite,CometArrayExpressionSuite,CometMapExpressionSuiteandCometNativeShuffleSuiteall pass.