Skip to content

feat: give each default Spark-to-Arrow conversion its own spark.comet.convert config - #6602

Merged
andygrove merged 1 commit into
apache:mainfrom
andygrove:feat/per-case-spark-to-arrow-configs
Oct 4, 2026
Merged

andygrove merged 1 commit into
apache:mainfrom
andygrove:feat/per-case-spark-to-arrow-configs

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6593.

Rationale for this change

spark.comet.sparkToColumnar.enabled and spark.comet.sparkToColumnar.supportedOperatorList decide which Spark leaf operators Comet converts to Arrow. The list holds Spark physical class names, and those change between Spark versions. Spark 4.1 plans a query without a FROM clause as a OneRowRelationExec. Before 4.1 it is an RDDScanExec, so the default OneRowRelation entry matches nothing there and the RDDScan entry converts it instead. A single switch also covers several cases with different trade-offs, so no case can have its own default or its own documentation. #6593 has the details. #832 gave Parquet, JSON and CSV their own spark.comet.convert configs for the same reason, and this PR does the same for the four operators the list names by default.

What changes are included in this PR?

It adds four configs, all off by default:

List entry Config
Range spark.comet.convert.range.enabled
InMemoryTableScan spark.comet.convert.inMemoryCache.enabled
RDDScan spark.comet.convert.rdd.enabled
OneRowRelation spark.comet.convert.oneRowRelation.enabled

CometExecRule checks each operator's own config first. A new shim, ShimCometOneRowRelation, recognizes a OneRowRelation on every Spark version. So before 4.1, the one-row config converts it and the RDD config does not.

No default changes, so nothing is converted unless a config is set, as before. The list now defaults to empty. It remains the way to convert other leaf operators, such as the BatchScan of a Data Source V2 connector. The old settings keep working:

  • When the list is not set, spark.comet.sparkToColumnar.enabled=true still converts the four operators above.
  • A list that names one of them still converts it.

Either way, Comet logs a warning once per JVM that names the config to use instead. The config docs, and the migration guide in a new 1.2.0 section, mark that use as deprecated.

Range gets its own config rather than becoming the fallback inside spark.comet.exec.range.enabled, which #6593 lists as a possible follow-up. Today the switch converts a range even with spark.comet.exec.range.enabled off, and every Comet suite relies on that through CometTestBase. Folding it in would change which operator those tests run.

Test changes:

  • CometTestBase now sets the four configs instead of the switch, which gives the same conversions on every Spark version.
  • Suites that turned the switch off, or narrowed the list to one operator, now set the matching configs.
  • CometInMemoryCacheSuite builds its own conf. Where it used to set the switch, a small helper now turns on the four configs.
  • CometTaskMetricsSuite still uses the list for LocalTableScan, which has no config of its own.

How are these changes tested?

There are three new tests in CometExecRuleSuite:

  • Each config converts its own operator and none of the others. This includes a query without a FROM clause, which Spark plans as an RDDScanExec before 4.1.
  • The old settings still convert what they converted before. The switch without a list converts all four operators, a list replaces them, and an operator without a config of its own is converted only when the list names it.
  • The deprecation warning appears once per config. It doesn't appear when the operator's own config is on, or when the list names an operator that has no config of its own.

I checked that each test fails when its part of the change is reverted:

  • Without the one-row case, the RDD config converts SELECT 1 on Spark 3.5.
  • Without the default for an unset list, the switch alone converts nothing.
  • Without the once-per-config guard, the warning appears twice.

Local runs:

  • Spark 4.1: CometExecRuleSuite, CometExecSuite, CometInMemoryCacheSuite, CometInMemoryCachePruningSuite, CometInMemoryCacheKryoSuite, CometArrayExpressionSuite, CometMapExpressionSuite, CometRangeExecSuite, CometEmptyRelationExecSuite, CometEmptyRelationParquetWriterSuite and CometTaskMetricsSuite pass.
  • Spark 3.5: CometExecRuleSuite, CometRangeExecSuite and CometInMemoryCachePruningSuite pass.
  • Spark 4.0: CometExecRuleSuite passes.
  • Semantic scalafix, scalastyle and spotless on 3.5, the strict-warnings compile on 3.5, the Spark 4.0 test compile and prettier are all clean.

I did not run the 3.4 and 4.2 profiles, which use the same shims as 3.5 and 4.1, or the plan stability suites.

….convert config

Add spark.comet.convert.range.enabled, spark.comet.convert.inMemoryCache.enabled,
spark.comet.convert.rdd.enabled and spark.comet.convert.oneRowRelation.enabled for the four
operators that spark.comet.sparkToColumnar.supportedOperatorList named by default. All are off
by default.

The list now defaults to empty and stays the way to convert other leaf operators. Using
spark.comet.sparkToColumnar.enabled and the list for the four operators above is deprecated: it
keeps working as before and logs a one-time warning naming the new config. The OneRowRelation
config matches both the RDDScanExec Spark plans before 4.1 and the OneRowRelationExec it plans
from 4.1.

CometTestBase and the suites that set the old keys move to the new ones.
@andygrove andygrove added the enhancement New feature or request label Oct 4, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Spark-to-Arrow conversion used one switch and physical operator names whose meaning varied across Spark versions.
  • Design approach: Add four disabled-by-default conversion configs while preserving the legacy switch and explicit operator list.
  • Correctness / compatibility analysis: The one-row shims match Spark sources across supported 3.4–4.2 profiles. An isolated probe using real Spark 4.1.3 plans passed 1,280 routing combinations. No introduced P1/P2 issues found within this review.
  • Key design decisions: Version differences remain in shims. The existing schema gate and native-operator precedence remain intact. The shared helper keeps the implementation small, and deprecation warnings are bounded to once per config per JVM.
  • Implementation sketch: Extend CometConf, add operator-specific routing in CometExecRule, and update test configuration and migration documentation. Added checks run during planning. No execution-path performance regression was identified.
  • Behavioral changes worth calling out: Compared with branch-1.1, the new opt-ins, empty list default and deprecation warnings are intentional and documented. An unset legacy list retains its effective conversions. Before Spark 4.1, the new one-row config distinguishes no-FROM queries from ordinary RDD inputs.
  • Suggested improvements: None meeting the P1/P2 reporting threshold.

Reviewed full SHA a2e626107e0f681f64d713d2049f6bdfec651e58 against supplied base 3bc2faa934f04799686705b952d59ed32a707dda. Inspected the tip differences and reconciled them with merge base fef94f6cd78b18151dff57b7a936798385356de5. The complete 22-file PR diff matches GitHub’s file list. Apparent Arrow/codegen reversions came from two newer base-only commits, not this PR. No existing reviews, issue comments, inline comments or threads contained unresolved concerns.

Skills: review-comet-pr, with expression, FFI and memory review siblings consulted while investigating and attributing the tip differences.

Exact-head CI at 2026-10-04 16:31 UTC: 22 successful checks, 30 skipped, six running, no reported failures. Spark 4.1’s four Comet test shards, Rust tests and TPC-DS remained running. Spark SQL and Iceberg suites were skipped.

Validation limits: The local routing probe used extracted production methods and a primitive-schema fixture, not the full Comet execution pipeline. No full native/JVM build or Spark SQL/Iceberg matrix was run locally. No project code or GitHub state was changed.

@andygrove
andygrove added this pull request to the merge queue Oct 4, 2026
Merged via the queue into apache:main with commit bdbc460 Oct 4, 2026
60 checks passed
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 5, 2026
…e of main

Since apache#6602, CometTestBase turns on the conversion of RDD scans with
spark.comet.convert.rdd.enabled rather than spark.comet.sparkToColumnar.enabled,
so setting the old key to false no longer kept the RDD scans in
CometShuffleInputConversionSuite on Spark. The scans were converted before
the shuffle input conversion saw them, and six tests failed. The suite now
turns off the conversions CometTestBase turns on, and the struct test and
the benchmark use the RDD scan config.
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 5, 2026
…park's cache scan

spark.comet.sparkToColumnar.enabled is deprecated for in-memory cached tables
since apache#6602, which gave the conversion its own config.
rich7420 pushed a commit to rich7420/datafusion-comet that referenced this pull request Oct 6, 2026
…che#6662)

Add spark.comet.convert.rowDataSource.enabled, off by default, which converts the output of
RowDataSourceScanExec to Arrow. Spark plans that scan for Data Source V1 relations that are not
file-based, such as JDBC tables. Until now the only way to convert it was to name
RowDataSourceScan in spark.comet.sparkToColumnar.supportedOperatorList.

That list entry still converts it, but as with the operators that got their own
spark.comet.convert config in apache#6602, it is now deprecated and logs a one-time warning naming the
new config. The list never named RowDataSourceScan by default, so the switch alone still does not
convert it.
dwsmith1983 pushed a commit to dwsmith1983/datafusion-comet that referenced this pull request Oct 7, 2026
…che#6607)

* feat: use native shuffle for row input by converting it to Arrow

When a shuffle's child is a Spark row-based operator, Comet uses its JVM
columnar shuffle. With the new spark.comet.convert.shuffleInput.enabled
(off by default), CometExecRule instead puts CometSparkToColumnarExec over
the child and uses native shuffle, where native shuffle supports the
partitioning and the conversion supports the columns. A shuffle that hashes
a decimal wider than 18 digits stays on the JVM columnar shuffle, because
native shuffle does not hash those as Spark does.

As for the typed Dataset conversion, the child's subtree gets its columnar
transitions from Spark's own rule first, since Spark inserts none below a
RowToColumnarTransition. The revert of a shuffle between two Spark
aggregates covers the new shape too.

* test: run CometShuffleInputConversionSuite in the PR builds

* fix: keep a shuffle that hashes a string on the JVM columnar shuffle

Spark's partitioner hashes a string's bytes as they are. Once
CometSparkToColumnarExec has converted the rows, the import into native
replaces invalid UTF-8 before native shuffle hashes the string, so such a
key can go to a different partition than Spark's partitioner sends it to.
A join whose other input stays on the JVM columnar shuffle then loses the
matches for that key. The conversion now leaves a shuffle that hashes a
string on the JVM columnar shuffle, as it does one that hashes a wide
decimal.

* fix: have native shuffle read a converted child's Arrow stream

A native shuffle over a CometSparkToColumnarExec wrapped the child's
executeColumnar batches in a ColumnarBatchArrowReader, which closes each
batch once native has it. Those batches share the vectors that the
conversion reuses for the next batch, and closing a struct vector drops
its children, so the next batch failed to import. Native shuffle now reads
such a child's Arrow stream, as a native operator does, which also fixes
the same failure with spark.comet.sparkToColumnar.enabled.

* fix: leave calendar intervals out of the shuffle input conversion

Arrow holds the time part of an interval in nanoseconds, so
IntervalMonthDayNanoWriter overflows on a calendar interval with more
microseconds than that can hold, which Spark accepts. The JVM columnar
shuffle leaves calendar intervals to Spark's shuffle, and the conversion
now does too.

* test: turn off RDD scan conversion with its own config after the merge of main

Since apache#6602, CometTestBase turns on the conversion of RDD scans with
spark.comet.convert.rdd.enabled rather than spark.comet.sparkToColumnar.enabled,
so setting the old key to false no longer kept the RDD scans in
CometShuffleInputConversionSuite on Spark. The scans were converted before
the shuffle input conversion saw them, and six tests failed. The suite now
turns off the conversions CometTestBase turns on, and the struct test and
the benchmark use the RDD scan config.

* fix: keep a shuffle whose key is computed from a string on the JVM columnar shuffle

The string guard looked only at the type of each hash key, so a key that
is a number computed from a string, such as hash(s), passed it. Native
shuffle evaluates such a key after the import into native has replaced
invalid UTF-8 in the string, so it can still send a row to a different
partition than Spark's partitioner does, and a join whose other input
stays on the JVM columnar shuffle loses the matches. The guard now looks
for a string anywhere in the key expression. A wide decimal still matters
only as the type of the key itself, because Comet's hash of a wide
decimal already falls back to Spark.

The string join test also joins on the hash of the strings, and uses two
keys instead of four. Both shuffles write the strings with invalid UTF-8
replaced, so the keys stay apart only while each key expression sends
them to different partitions, and two of the four hashes share one.

* fix: put a converted shuffle back on the JVM columnar shuffle when its stage is reverted

RevertNativeForTransitionHeavyStages strips CometSparkToColumnarExec with the rest of a
reverted stage, which left a native shuffle reading the batches of Spark's RowToColumnarExec,
and the task failed casting them to CometVector. The shuffle now goes back to the JVM columnar
shuffle that the conversion replaced, which reads the stage's rows.

* refactor: convert only the input of a shuffle that would use the JVM columnar shuffle

convertsInputForNativeShuffle now starts from shuffleSupported, whose checks it restated, and
which also rules out Celeborn and calendar intervals. With spark.comet.shuffle.mode=native, a
Spark child keeps Spark's shuffle, as the config doc says.

* perf: sample a conversion's input rows for range partitioning

The range partitioner reads only the sort keys, so it samples the rows the conversion reads
rather than converting all of their columns to Arrow first. The benchmark gets a
range-partitioned case.

* test: check where native shuffle puts each key type the conversion admits

Checks the partition of every row and a join against the JVM columnar shuffle for each hash key
type the conversion admits, with NULL and boundary values. Adds array<string> and
map<string,string> columns with NULLs to the test over several batches, and plan assertions to
the test of aggregates whose partial aggregate runs on Spark.

* docs: describe the shuffle input conversion in the contributor guide

Also corrects the revert's log message in the tuning guide.

* test: count a range-partitioned conversion's rows once, and range-partition every column in the benchmark

Before the sampling change, the range partitioner's sampling job ran through the conversion and
added its rows to the conversion's metrics, so a range-partitioned shuffle reported twice the
rows it wrote.

* docs: say that the conversion gains with hash partitioning and comes out even with range partitioning

Also rewords the contributor guide's rule for when native shuffle takes a Spark child.

* test: correct why the shuffle between the two sort aggregates stays native

Since apache#4565, the revert of a shuffle between two Spark aggregates takes sort aggregates too, so
the reason the test gave, that the revert covers hash aggregates only, no longer held. The
shuffle stays native because the final sort aggregate reads it through a sort.

* docs: name sort aggregates in the tuning guide's shuffle revert section

Since apache#4565, the revert of a shuffle between two Spark aggregates takes SortAggregateExec too, but
only where the final aggregate reads the shuffle directly. One with grouping keys reads it through
a sort, so that shuffle is not reverted. The last sentence also no longer calls the shuffle that
disabling the revert keeps a columnar shuffle, since with the shuffle input conversion it can be a
native one.

* docs: say that the shuffle revert takes any Spark aggregate

Since apache#4565, the revert of a shuffle between two Spark aggregates matches any BaseAggregateExec,
not only hash aggregates, so the config doc and the scaladoc of revertRedundantColumnarShuffle
no longer name hash aggregates alone. The scaladoc also says why the shuffle below a sort
aggregate with grouping keys stays.

* test: check the Celeborn exclusion of the shuffle input conversion in every shuffle mode

In native mode the JVM columnar shuffle is off, so the conversion stays off
with or without the Celeborn checks. Run the test in auto and jvm too, where
Celeborn is what keeps the Spark child on Spark's shuffle, and check that the
fallback reason names Celeborn, which in native mode only the Celeborn check
in shuffleSupported gives.

* fix: keep native shuffle over a stacked bridge when reverting a transition-heavy stage
parthchandra pushed a commit to parthchandra/datafusion-comet that referenced this pull request Oct 9, 2026
* perf: project cached batches by buffer selection, prune on collated strings

Follow-up to apache#5051, applying items from apache#5487.

Replace the per-column Arrow IPC stream layout of `CometCachedBatch` with a
single encapsulated IPC record batch message per cached batch, carrying no
Schema message and no end-of-stream marker. The reader rebuilds the schema
from the cached relation's attributes, so a wide relation no longer repeats
the same schema bytes once per cached batch.

Compression moves from a whole-payload Spark codec to Arrow's per-buffer IPC
compression. That is what makes projection cheap: the message metadata records
every buffer's offset and length in the body, so `CachedBatchIpc.readProjected`
copies out only the byte ranges of the columns a scan selected and decompresses
just those. This subsumes the separate "drop the schema message" item, since
there is no longer a per-column stream to frame.

Dictionary-encoded columns are decoded before being stored: a payload with no
schema message cannot describe a dictionary encoding.

The codec defaults to zstd, and lz4 is deliberately not offered. Arrow's lz4 is
commons-compress's pure-Java implementation, unrelated to the JNI-accelerated
lz4-java behind `spark.io.compression.codec`. Over a 200k-row six-column
relation it measured 205s to write against 347ms for zstd, while also producing
larger output, so no workload prefers it. zstd also beats storing batches
uncompressed on both axes (347ms and 2 MiB against 1743ms and 13 MiB), because
the bytes it saves cost more to copy and store than compressing them costs.

Decompression is done here rather than left to `VectorLoader`, which leaks:
`VectorLoader.loadBuffers` collects a field's decompressed buffers into a local
list and releases them only after the whole field loads, so a buffer that fails
to decompress strands every buffer of that field decompressed before it. A
string column reaches this, its offsets buffer decompressing before its data
buffer throws.

Also track statistics bounds for collated string columns, comparing with the
collation's own ordering through a new `CometTypeShim.compareStrings`. Matching
the bare `StringType` object excluded collated columns, which then got null
bounds and no pruning.

Benchmark over a 5M-row six-column relation, keeping the cached scan native
against falling back to a Spark cache scan and converting: 1.3x on a repeated
scan, 1.3x on a narrow projection and 2.3x on a full projection.

* refactor: use Spark's interpreted ordering for bounds, hoist projection layout

Cleanup pass over the cache format change. No behaviour change.

Drop the `compareStrings` shim in favour of `TypeUtils.getInterpretedOrdering`.
That method is public with the same signature on every supported Spark version,
and on Spark 4 it resolves a `StringType` through
`CollationFactory.fetchCollation(collationId).comparator` -- the comparison the
shim was reaching for. So the collation awareness comes from Spark itself and
the shim, its Spark 3.x stub and the hand-rolled per-type `compare` all go.
The ordering is now resolved once per column per partition rather than being
re-dispatched on the `DataType` twice per row.

Build the projection's index layout once per partition instead of per batch.
The node, buffer and variadic index arithmetic is a pure function of the cached
schema and the selected columns, but it walks every field of the relation, so
recomputing it per batch made the bookkeeping O(total columns) against O(selected
columns) of useful work -- worst in the wide-relation, narrow-projection case the
format exists for. `CachedBatchIpc.Projection` now holds that layout and the
projected schema, and owns the whole decode; `ProjectedBatch` is left with
ownership only. This also puts the projected schema next to the code that packs
buffers in the same order, an invariant that previously spanned two files
unstated.

Smaller cleanups: use Arrow's `DataSizeRoundingUtil.roundUpTo8Multiple` rather
than open-coding IPC body alignment; size the serialization buffer from the
record batch's known body length instead of growing from 32 bytes; resolve
decompressors once instead of per batch; share the dictionary lookup guard
between `Utils.combineDictionaryProviders` and the cache writer; read the codec
config through one helper carrying the driver-vs-executor rationale; and collapse
the duplicated compressed-buffer predicate and scramble loop in the test helper.

Corrects two `Utils` scaladocs that still described the per-column stream format
this change replaced. Benchmark and codec figures in the docs re-measured against
the current code.

* fix: relocate the arrow-compression service file when shading

arrow-compression ships
META-INF/services/org.apache.arrow.vector.compression.CompressionCodec$Factory.
The shade plugin copies it verbatim without a ServicesResourceTransformer, so
the jar declared a provider for Spark's own unshaded Arrow interface while
naming a class that exists here only under the relocated package. Every
ServiceLoader lookup Spark's Arrow made then failed with a
ServiceConfigurationError, which took CompressionCodec.Factory's static
initializer down with it and broke unrelated Arrow IPC reads, including
mapInArrow.

Add ServicesResourceTransformer so the service file name and its contents are
both relocated. arrow-compression is the only bundled artifact that ships one.

Also drop an unused NonFatal import that scalafix flagged.

* test: drop the cache leak test that depends on zstd corruption detection

"releases its vectors when a column fails part way through" zeroed the last
16 bytes of a compressed buffer and required the read to fail. Whether that
fails is a property of the zstd runtime, not of Comet: the cached payload is
byte-identical across Spark versions, but Comet takes zstd-jni from Spark
rather than from arrow-compression, and 1.5.5 (Spark 3.4, 3.5) decodes that
frame while 1.5.7 (Spark 4.x) reports it corrupt. So the test passed on 4.x
and failed on 3.4 and 3.5.

The scenario it claimed to cover is also unreachable: CachedBatchIpc
decompresses every selected buffer before VectorLoader runs, so no content
corruption can fail part way through the load. The two remaining leak tests
corrupt a frame from its header onwards, which every zstd release rejects,
and already cover a failure at a column's first buffer and a failure after an
earlier buffer of the same column decoded.

Records the constraint on scramble so a future test does not reach for a
tail-only corruption again, and drops the now unused truncateColumn helper
and the dictionary fixture's payload argument.

* feat: enable Comet's in-memory cache by default

Flip spark.comet.exec.inMemoryCache.enabled to true so cached tables are
stored and scanned in Comet's Arrow format without an opt-in.

CometDriverPlugin.maybeSetCacheSerializer read the config out of SparkConf
with a hardcoded false default, so flipping the ConfigEntry alone would have
left the serializer uninstalled unless the user set the key explicitly. It
now falls back to the entry's own default, matching how the plugin reads
spark.comet.metrics.enabled.

Stacked on apache#5543.

* test: cover nested columns in the cached-batch projection tests and benchmark

Addresses review feedback asking whether nested data should be tested and
benchmarked.

Nested columns were already round-tripped, but only under a full projection,
which cannot see the part of the format that is nontrivial for them. A flat
column always owns one field node and two or three buffers; a nested one owns a
run as long as its subtree, and selecting every column covers the whole
sequence however it is partitioned. So the buffer-span arithmetic was only
exercised in the one shape where getting it wrong does not show.

Adds two tests over a six-column relation whose middle four columns are a
struct, an array, a map and a struct wrapping an array:

- Each column takes its turn as the sole projection while the other five are
  corrupted, so a run computed short or long is caught by reaching into a
  corrupted neighbour.
- Values are compared against the uncached query across single-column,
  paired and out-of-order projections. Row counts cannot catch a window that
  is misaligned but still decompresses, and out-of-order is the case a full
  projection cannot stand in for.

The per-column statistics test now runs over the nested relation too, since a
nested column's recorded size is the sum of its whole subtree.

Both new tests fail if fieldNodeCount stops recursing into children.

In the benchmark, adds the three projection widths over a relation of struct
columns, and asserts the width each case claims. That assertion caught the
existing "full projection (6 of 6 columns)" case reading three: count() over a
non-nullable column is rewritten to count(1) by NullPropagation, which prunes
the column out of the scan, and only k, s1 and s2 were nullable -- and those
only incidentally, because Remainder can divide by zero. Every column of both
relations is now nullable so count(c) genuinely reads c, and the documented
numbers are regenerated.

Array and map columns are left out of the benchmark deliberately: the baseline
arm needs Spark's cache scan to bridge into Comet operators, and
CometSparkToColumnarExec declines ArrayType and MapType, so for those the arm
does not exist and the two cases stop measuring the same boundary. The docs say
so rather than leaving it to be rediscovered.

* fix: drop a redundant string interpolator flagged by scalafix RedundantSyntax

* review: check the cached layout, and address the rest of the review

Reader-side: a cached payload carries no schema, so `Projection` derived
every node and buffer window from `Utils.toArrowSchema(cacheAttributes)`
with nothing checking the writer had produced that layout. `load` now
compares `nodesLength()`/`buffersLength()` against the totals
`selectedRange` already computes, before any unchecked `batch.buffers(j)`.

Writer-side: `isArrowBacked` accepts a `FixedSizeBinaryVector` for a
`BinaryType` column, which is two buffers where the reader rebuilds three,
and it answers for the top-level vector only -- so a struct of large
strings passes it and is stored with 64-bit offsets. `matchesReaderLayout`
compares the batch's Arrow types against the reader's recursively, and a
batch that disagrees takes the conversion path instead. A dictionary
column's field carries the index type, so the dictionary's field is what
is compared.

Also: an unrecognized body-compression byte is rejected rather than read
as plain bytes, `fieldVariadicCount` and the variadic plumbing are gone
(the length check covers view vectors, which the counts would not have),
`columnSizes` no longer re-walks each column's subtree, the write codec is
a case class rather than a bare tuple, the per-partition `Projection` is
lazy so a row-count-only read never builds it, `hydrateDictionaries` is
`decodeDictionaries`, `Projection` takes an `IndexedSeq`, and the stale
`readProjected` links and some over-long comments are fixed.

Tests: the two projection tests become one parameterized over both
relations, caching once and restoring the payload between columns instead
of re-caching; the two leak tests become one with two corruption points.
New tests cover the reader's layout check and the writer declining a
fixed-size-binary batch.

* review: own the write-side compression buffers, fix the activation example

Compressing through VectorUnloader leaks on the failure path: appendNodes
retains each input buffer and accumulates the compressed ones into a list
local to getRecordBatch, so a buffer that fails to compress strands that
retain and leaves every buffer compressed before it reachable from nothing.
Closing the input batch afterwards undoes neither. Unload plain and compress
in CachedBatchIpc.compressed instead, mirroring what decompressed already
does on the read side, so every allocation stays reachable from an error
path that owns it.

The docs enabled the cache with spark.conf.set, which cannot work: the
driver plugin picks spark.sql.cache.serializer while the SparkContext is
initializing. Show it as a startup --conf.

Also drops a redundant s interpolator that the scalafix lint rejected.

* fix: write cached batches to the schema width, not the batch width

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.

* fix: report a cached relation's decoded size to the planner, not its compressed size

Comet's cache format reported each batch's sizeInBytes as its stored,
compressed payload. Spark's planner reads the sum of those as the size of a
materialized cached relation, for the broadcast threshold and the shuffled
hash join build side among others, and both of Spark's own cache formats
report decoded sizes there. With the default zstd codec a cached relation
could look several times smaller than in Spark's format and be broadcast
where Spark would shuffle it.

Record each column's decoded Arrow size, measured before compression as
Spark's ArrowCachedBatchSerializer does, and let CometCachedBatch inherit
sizeInBytes from SimpleMetricsCachedBatch as the sum of those.

* fix: keep Spark's cache scan for a relation whose cached plan records observed metrics

Spark collects Dataset.observe metrics after a query with
CollectMetricsExec.collect, which reaches the metrics recorded inside a cached
plan only through an InMemoryTableScanExec over it. With Comet's native cache
scan in its place those metrics were lost: QueryExecution.observedMetrics came
back empty, Observation.get returned an empty map on Spark 3.5 and later, and
on Spark 3.4 it never returned.

Keep Spark's scan for such a relation, with a fallback reason. The data stays
in Comet's format and is read through the existing fallback path. The check
walks the cached plan the way CollectMetricsExec.collect does, through
subqueries, adaptive plans and nested caches.

* test: install Comet's cache serializer in the Spark SQL test sessions

The Spark SQL diffs turn Comet on through SharedSparkSession and TestHive
rather than by loading CometPlugin, and the plugin is what installs Comet's
cache serializer. So even with spark.comet.exec.inMemoryCache.enabled on by
default, every Spark SQL suite cached in Spark's own format and none of them
exercised Comet's.

Set spark.sql.cache.serializer in both places when Comet is enabled, as the
plugin would, and adapt the tests that assume Spark's cache scan:

- plan checks that only need a cache scan to be present also accept
  CometInMemoryTableScanExec, reading its originalPlan where the test uses
  the scan's relation;
- PartitionBatchPruningSuite checks its answers as before, but reads the
  test-only accumulators only from Spark's scan, which Comet's lacks;
- CacheTableInKryoSuite registers Comet's classes with
  CometKryoRegistrator, as Comet asks of any application that sets
  spark.kryo.registrationRequired;
- tests of Spark internals that Comet's scan replaces are tagged
  IgnoreComet: exact cache size estimates, the union's columnar support,
  subquery reuse through the scan's predicates, and AQE coalescing under a
  union (apache#6454).

Each diff was regenerated from a clone of its Spark tag with the existing
diff applied, after checking that it round-tripped unchanged.

* fix: coalesce shuffle partitions under Comet unions the way Spark does

Spark's CoalesceShufflePartitions coalesces each child of a UnionExec as its
own group, and from Spark 4.0 each child of a CartesianProductExec,
BroadcastHashJoinExec or BroadcastNestedLoopJoinExec too. It matches those
classes, so a CometUnionExec fell through to the case that coalesces only when
every leaf below it is an exchange stage. A union with a scan or a table-cache
stage in one branch kept every partition of the shuffles in the others.

Add CometCoalesceShufflePartitions, a query-stage optimizer rule that runs
after Spark's. For a Comet operator whose Spark original is one of those, and
whose shuffle stages no AQE rule has touched, it rebuilds the Spark operators
over the Comet children, runs Spark's own rule over that, and swaps the Comet
operators back in. It extends AQEShuffleReadRule, so AQE gates and validates
it the way it does Spark's rule. Spark 3.4 has no hook for query-stage
optimizer rules, so the fix applies from Spark 3.5.

Closes apache#6454.

* test: fix three Spark SQL test adaptations for Comet's cache format

- CachedBatchSerializerNoUnwrapSuite sets its own cache serializer, but
  Spark keeps the first one it loads for the rest of the JVM, so the
  suite ran with Comet's. Clear it before and after the suite, as
  CachedBatchSerializerSuite does (4.1.3, 4.2.0).
- The SPARK-19993 helper counted the subqueries in a Comet cache scan's
  original plan, which subquery reuse never reaches, on top of the same
  subqueries in the filter above the scan.
- SPARK-35332's AQE check reads the cached plan from the SQL UI's plan
  tree, which shows it only under Spark's own cache scan (apache#6463). Skip
  it under Comet (4.0.4 and later).

* fix: keep Spark's cache format when Kryo would reject Comet's or Comet disables itself

CometDriverPlugin installed ArrowCachedBatchSerializer whenever Comet, its
native execution and the cache config were enabled at startup. Two more
startup settings decide whether that format can work:

- With spark.kryo.registrationRequired=true and no CometKryoRegistrator,
  Kryo rejects CometCachedBatch, so caching failed with "Class is not
  registered" as soon as Spark serialized a cached block, where Spark's
  own format works. The plugin only warned. It now keeps Spark's format,
  and the warning says so.
- With Comet shuffle enabled but neither of Comet's shuffle managers
  configured, Comet disables itself, so every cache was stored in Comet's
  format with only Spark operators to read it. The plugin now keeps
  Spark's format there too.

The plugin's boolean config reads also honor a deprecated alternative
key, such as spark.comet.exec.shuffle.enabled, as a session does.

* docs: describe the in-memory cache default in the upgrade guide

Add a 1.2.0 section to the upgrade guide for the new default of
spark.comet.exec.inMemoryCache.enabled: what changes, how to keep Spark's
cache format, and when Comet keeps it without being asked. The operator
and compatibility pages still called the feature disabled by default,
and the cache guide still said an application that started with the
default kept Spark's format.

* test: benchmark the in-memory cache against Spark's format with AQE on

The published comparisons either read the same Comet-written cache both
ways, with spark.comet.sparkToColumnar.enabled turned on, or read both
formats with Comet off. Neither is what an application gets from the
feature's default. Add cases with Comet and AQE on and Comet's other
settings at their defaults, reading a relation cached in Spark's format
and in Comet's, with Comet operators above the cache scan and with a
Spark operator above it.

* feat: record why Spark scans a relation cached in Comet's format when Comet is off

spark.sql.cache.serializer is static, so a relation cached in Comet's
format keeps it for the life of the application. A session that turns
Comet or its native execution off then reads it with Spark's
InMemoryTableScanExec, which is slower than reading Spark's own format
(apache#5485). CometExecRule records a fallback reason for the other ways Spark
ends up scanning a cache, but it returned before looking at the plan in
these two cases, so nothing explained them.

Record a reason on each such scan, naming the cause and the startup
setting that keeps caches in Spark's format.

* docs: say an application keeps the cache format it started with

The sentence described a session started with the default, which keeps
Spark's format only while the default is off.

* test: format the in-memory cache benchmark

* fix: bind the fallback reason result for the strict Scala warnings build

-Ywarn-value-discard rejects the node that withFallbackReason returns,
discarded as the last expression of a case.

* docs: compare the in-memory cache with Spark's format as an application gets it

Add the benchmark's adaptive results to the cache guide: with Comet and
AQE on and Comet's other settings at their defaults, Comet's format is as
fast or faster than Spark's in every shape but a read of three long
columns, which zstd decompression makes about 10% slower. That includes
a Spark operator above the native scan.

The guide said that, without the feature, a CometSparkColumnarToColumnar
converts each batch of Spark's cache scan for Comet. That takes
spark.comet.sparkToColumnar.enabled, which is off by default, so by
default the operators above Spark's scan run on Spark. The limitation
that Spark operators read Comet's format slowly applies to Spark's own
cache scan, not to Spark operators above Comet's.

* fix: count Kryo registrations however the application made them

The gate kept Spark's cache format whenever spark.kryo.registrator did
not list CometKryoRegistrator, even where the application had registered
Comet's cached batch another way, such as spark.kryo.classesToRegister.
Comet's format works there, while Spark's does not before 4.1, which
registers DefaultCachedBatch itself only from then on, so the gate broke
caches that worked.

Ask a Kryo instance built from the application's conf which of Comet's
classes it has registered, and keep Spark's format only if Comet's
cached batch is not among them. If Kryo cannot be built, fall back to
looking for CometKryoRegistrator in spark.kryo.registrator. The startup
warning uses the same check and names the missing classes.

* docs: describe the Kryo condition by registration in the upgrade guide

The plugin now keeps Spark's cache format when Kryo has not registered
Comet's cached batch, however the application registers classes, rather
than when spark.kryo.registrator does not list CometKryoRegistrator.

* docs: name spark.comet.convert.inMemoryCache.enabled for converting Spark's cache scan

spark.comet.sparkToColumnar.enabled is deprecated for in-memory cached tables
since apache#6602, which gave the conversion its own config.

* test: accept Comet's cache scan in BroadcastJoinSuite's cached-relation tests

Since apache#6415, suites that build their own SparkSession run Comet, and with the
cache on by default BroadcastJoinSuite's cached relations are read by
CometInMemoryTableScanExec. SPARK-23214 and SPARK-37742 count Spark's cache
scan and broadcast hash join, so count Comet's as well.

* fix: coalesce with the whole stage in view and keep partitionings above Comet unions

Address review of the CometCoalesceShufflePartitions rule:

- Stand in only for Comet operators whose output partitioning is unknown.
  From Spark 4.1 a union of co-partitioned children reports their
  partitioning, and an aggregate above reads it without a shuffle. Coalescing
  one branch made Comet's aggregate, which states no required distribution,
  return each key twice.
- Hand Spark's rule the whole stage rather than the subtree below each Comet
  operator, so it sees a Cartesian product or nested loop join above a union
  and uses its smaller target size there, and divides the minimum partition
  count over every group. Leave a stage that already has a read over a
  shuffle, which Spark's rule would coalesce again.
- Copy the Spark original when withNewChildren hands it back, so a union that
  AQE plans again over its stages is coalesced too.
- Drop the CartesianProductExec case, which has no Comet counterpart, and
  document the rule in CometSparkSessionExtensions.

* test: run Spark's SPARK-42101 union coalescing test with Comet again

The test failed with Comet's cache format because AQE did not coalesce the
shuffled branch of a union that Comet ran (apache#6454), which apache#6459 fixes.

* fix: coalesce Comet union branches in stages that already hold an AQE shuffle read

CometCoalesceShufflePartitions returned a stage unchanged when it held
any AQEShuffleReadExec, because Spark's rule asserts on reads that are
not skew splits. Below a Cartesian product, Spark's rule coalesces the
shuffle on the other side by itself first, so the read it leaves there
kept the shuffled branch of a Comet union at every partition.

Hide each read behind a leaf while Spark's rule runs, and put it back
afterwards. Spark's rule then leaves alone the shuffles it would
coalesce together with a read one, and coalesces the rest of the
stage.

* docs: say Spark's cache scan can read Comet's format more slowly

The fused reader from apache#5859 reads Comet's format about as fast as
Spark's own, and faster for the narrowest reads, so the upgrade note
should not say Spark's scan is always slower.

* test: leave Comet's cache scan out of the DPP suite's subquery check

Since apache#6577, CometInMemoryTableScanExec exposes Spark's own
InMemoryTableScanExec as its one subquery, so that Spark's UI draws the
cached plan below it. With the cache on by default, the DPP suite's
'filtering ratio policy fallback' caches its dimension table, and
checkPartitionPruningPredicate requires every subquery of an adaptive
plan to contain an AdaptiveSparkPlanExec, which Spark's scan does not.
Skip it there in every diff, since the query never runs it as a
subquery.

* feat: keep Comet's cache format opt-in on Spark 3.4

Spark 3.4 has no hook for query-stage optimizer rules, so
CometCoalesceShufflePartitions cannot run there, and AQE cannot
coalesce the shuffled branch of a Comet union whose other branch reads
a cache in Comet's format (apache#6454). spark.comet.exec.inMemoryCache.enabled
now defaults to true only from Spark 3.5, and the plugin, which follows
the entry's default, keeps Spark's format on 3.4 unless it is set.

The 3.4 Spark SQL diff goes back to main's: its sessions install
Comet's serializer only as the plugin would, which on 3.4 it no longer
does by default, so its cache test adaptations are not needed there.

* docs: rewrap the Spark 3.4 cache default sentences

* test: run Spark's SPARK-35332 AQE cache test with Comet again

The diffs skipped CachedTableSuite's 'SPARK-35332: Make cache plan
disable configs configurable - check AQE' on Spark 4.0 and later,
because SparkPlanInfo drew the cached plan only below Spark's own
InMemoryTableScanExec. Since apache#6577 it is drawn below
CometInMemoryTableScanExec too, so the test finds the cached plan.

The test's check that the cached plan was coalesced, which the diffs
already adapt, looked for AQEShuffleRead only below a columnar-to-row
transition, where it is when Comet's plan is cached in Spark's format.
In Comet's columnar format the cached plan ends in the read itself, as
in Spark, so the check accepts it there too. The test passes on Spark
4.0, 4.1 and 4.2 with either format, and with Comet off.

* docs: stop describing the in-memory cache as experimental

It is enabled by default from Spark 3.5. The operator table now marks
InMemoryTableScanExec as supported, and the plan-node table no longer says
the cache scan is disabled by default.

* test: skip only the format-dependent checks in the Spark SQL diffs

SPARK-36120, SPARK-33687, SPARK-22673 and SPARK-43376 now run with Comet's
cache format and skip only the size or subquery-reuse assertions that depend
on it. The union test in CoalesceShufflePartitionsSuite runs its coalescing
check again on Spark 3.5 and later, where CometCoalesceShufflePartitions
coalesces the shuffles on the union's join side.

From Spark 3.5, CometTestSettings also makes Comet's cache serializer the
default of the sql and hive test JVMs, next to the shuffle manager. Spark
keeps the first cache serializer it loads for the rest of the JVM, so a suite
that caches first through a session of its own left the suites after it in
Spark's format.

* test: time building the cache in Spark's format in the cache benchmark

runBuildBenchmark caches the same relation in Spark's format and in Comet's,
with each codec, from Comet operators and from Spark operators, and reports
the footprint of each. The cache guide has the results.

* test: skip only the columnar checks of SPARK-37371 under Comet

Comet runs the union of two cached relations as a CometUnion over its own
cache scans, so the checks on Spark's scans and union are skipped there. The
union with a local relation stays Spark's non-columnar UnionExec, and both
answers are checked, under Comet as well.

---------

Co-authored-by: test <a@b.c>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Iceberg enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Give each default Spark-to-Arrow conversion its own spark.comet.convert.* config

2 participants