fix: write cached batches to the schema width, not the batch width - #6090
Conversation
ArrowWriter.writeColumns drove its loop from the input ColumnarBatch's width while indexing the writer's fields, which come from the schema the batch is written under. That assumed every producer hands over a batch exactly as wide as the schema. Iceberg's vectorized reader does not. BatchDeleteFilter.filterBatch reads with the delete filter's requiredSchema, which carries _pos after the projected columns when a data file has position deletes, and trims the extras back only when the file also has equality deletes. A merge-on-read UPDATE writes position deletes and no equality deletes, so the extra column survives into the batch, and caching such a relation failed with ArrayIndexOutOfBoundsException inside the write loop. Drive the loop from the writer's fields instead, which writes exactly the columns the schema describes: the extras are trailing, the same prefix Iceberg keeps when it does trim. A batch narrower than the schema is a genuine contract violation and is now refused with a message naming both widths. Closes apache#6087.
sunchao
left a comment
There was a problem hiding this comment.
Correctness
The schema-width fix addresses the reported cache failure at e483d083aecc against 09b44ad6fa17. Previously, a one-field writer receiving [name, _pos] wrote name, then indexed nonexistent fields(1). The new loop writes the schema's leading columns and explicitly rejects a batch narrower than that schema before writing any column.
I traced the producer and reader contracts. Iceberg 1.11 builds its required schema by retaining the requested fields in order and appending delete-only fields. Its position-delete path can retain those extras. On the maintained Spark 3.5/4.0 branches, InMemoryRelation supplies cachedPlan.output alongside the plan's columnar batches, and the Arrow row writer also follows its fields' width. Comet builds the conversion schema from those attributes, including their types and nullability. Cache statistics use the attribute width, and cache reads resolve requested columns by expression ID. Ordinary partition columns that belong to the declared output are preserved. This is a prefix-width fix, with the existing type/null writers and slice offsets unchanged.
Zero-column row counts remain explicit. Empty input, null handling and split batches retain their existing paths. The new narrow-batch exception is covered by the converter's existing root cleanup, and the new test verifies zero outstanding allocator bytes. I found no new or remaining P1/P2 issue.
One scope clarification: this exact head uses Utils.isArrowBacked and Utils.serializeBatchColumns. It does not contain the CachedBatchIpc.matchesReaderLayout guard mentioned in the description. The Arrow-backed branch predates this PR and is unchanged. This fix repairs the non-Arrow conversion branch.
Validation
Both new regression tests passed in the Spark 4.1 execution job, which completed with 976 succeeded, zero failed, zero canceled and five ignored. Its actual checkout was 486b2aabe685, whose parents are the reviewed base and head and whose whole tree matches the head. Current checks show 24 successes and 14 skips, including successful Required Checks. The skipped Spark SQL and Iceberg integration jobs do not establish an end-to-end merge-on-read cache pass.
Local validation was source analysis, git diff --check, and a source-derived loop model that reproduces the base overrun and verifies the new prefix bound. The model is not a JVM or Arrow runtime test. Maintained Spark 3.5/4.0 sources were checked. Maintained 3.4/4.1 sources remain unavailable. No local full suite or benchmark was run.
Performance
The change adds one width check per conversion or slice and bounds the existing column loop by the schema. It adds no per-row work, replacement batch, projection array or extra copy. The existing typed bulk-copy and scalar fallback paths remain unchanged. The implementation avoids copying trailing connector columns, and there is no credible new material cost that calls for a benchmark for this correctness fix. No measured performance improvement is claimed.
Design
Putting the bound in ArrowWriter.writeColumns keeps both cache conversion and SparkColumnarArrowReader consistent with the row writer. The explicit lower-bound check distinguishes valid trailing producer fields from missing required fields. The fix leaves cache format, selection, ownership and exception cleanup with their existing owners. That is a small, appropriate change at the shared conversion boundary.
Abstraction & complexity
No new abstraction or configuration is introduced. The existing writer fields are already the source of schema width, so using their length avoids an additional connector-specific projection layer. The comment explains the otherwise surprising wider-batch contract, and the two tests directly cover acceptance of the requested prefix and rejection with cleanup. I have no further actionable simplification to propose.
Which issue does this PR close?
Closes #6087.
Rationale for this change
ArrowWriter.writeColumnsdrives its loop from the width of the inputColumnarBatchwhileindexing the writer's fields, which are built from the schema the batch is being written under:
That assumes every producer hands over a batch exactly as wide as the schema. Iceberg's vectorized
reader does not.
BaseBatchReader.BatchDeleteFilter.filterBatchreads withdeletes.requiredSchema(), which carries_posafter the projected columns when a data file hasposition deletes, and it only trims the extras back with
ColumnarBatchUtil.removeExtraColumnswhen the file also has equality deletes. A merge-on-read UPDATE writes position deletes and no
equality deletes, so the extra column survives into the batch.
Caching such a relation crashes.
InMemoryRelationpassescachedPlan.outputas the cache schema,and because
supportsColumnarInputis true Spark strips theColumnarToRowabove the cached plan,so the serializer receives the scan's batches unaltered. Caching
SELECT name FROM toverid INT, name STRINGthen gives the writer a 2-column batch and 1 field:Comet's own row path already does the right thing:
ArrowWriter.write(row)loopsfields.length.So does Spark's
ArrowCachedBatchSerializer(SPARK-57268), which converts a non-Arrow batchthrough
batch.rowIterator()and the same schema-driven writer.writeColumnswas the one paththat trusted the batch instead of the schema.
What changes are included in this PR?
writeColumnsis now driven by the writer's fields rather than the input's width, so trailingcolumns the schema does not describe are ignored. Those are exactly the columns Iceberg itself
discards when it trims, since
removeExtraColumnskeeps the leadingexpectedSchemaprefix.The opposite mismatch is a real contract violation and stays fatal, but it now fails with a message
that names both widths instead of an
ArrayIndexOutOfBoundsExceptionfrom inside the write loop.The direct-write path needed no change:
CachedBatchIpc.matchesReaderLayoutalready requiresbatch.numCols() == readerFields.length, so a wider batch declines it and falls into theconversion path this PR fixes.
How are these changes tested?
Two tests in
CometArrowStreamSuite, which covers this boundary contract directly: a batch widerthan the schema is written as the schema describes it, with the trailing column ignored and the
projected values intact; a batch narrower than the schema is refused, with the allocator left at
zero so the failure does not leak the partially built root.