Repository navigation
feat: give each default Spark-to-Arrow conversion its own spark.comet.convert config - #6602
Conversation
….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.
sunchao
left a comment
There was a problem hiding this comment.
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 inCometExecRule, 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.
…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.
…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.
…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.
…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
* 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>
Which issue does this PR close?
Closes #6593.
Rationale for this change
spark.comet.sparkToColumnar.enabledandspark.comet.sparkToColumnar.supportedOperatorListdecide 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 aFROMclause as aOneRowRelationExec. Before 4.1 it is anRDDScanExec, so the defaultOneRowRelationentry matches nothing there and theRDDScanentry 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 ownspark.comet.convertconfigs 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:
Rangespark.comet.convert.range.enabledInMemoryTableScanspark.comet.convert.inMemoryCache.enabledRDDScanspark.comet.convert.rdd.enabledOneRowRelationspark.comet.convert.oneRowRelation.enabledCometExecRulechecks each operator's own config first. A new shim,ShimCometOneRowRelation, recognizes aOneRowRelationon 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
BatchScanof a Data Source V2 connector. The old settings keep working:spark.comet.sparkToColumnar.enabled=truestill converts the four operators above.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.
Rangegets its own config rather than becoming the fallback insidespark.comet.exec.range.enabled, which #6593 lists as a possible follow-up. Today the switch converts a range even withspark.comet.exec.range.enabledoff, and every Comet suite relies on that throughCometTestBase. Folding it in would change which operator those tests run.Test changes:
CometTestBasenow sets the four configs instead of the switch, which gives the same conversions on every Spark version.CometInMemoryCacheSuitebuilds its own conf. Where it used to set the switch, a small helper now turns on the four configs.CometTaskMetricsSuitestill uses the list forLocalTableScan, which has no config of its own.How are these changes tested?
There are three new tests in
CometExecRuleSuite:FROMclause, which Spark plans as anRDDScanExecbefore 4.1.I checked that each test fails when its part of the change is reverted:
SELECT 1on Spark 3.5.Local runs:
CometExecRuleSuite,CometExecSuite,CometInMemoryCacheSuite,CometInMemoryCachePruningSuite,CometInMemoryCacheKryoSuite,CometArrayExpressionSuite,CometMapExpressionSuite,CometRangeExecSuite,CometEmptyRelationExecSuite,CometEmptyRelationParquetWriterSuiteandCometTaskMetricsSuitepass.CometExecRuleSuite,CometRangeExecSuiteandCometInMemoryCachePruningSuitepass.CometExecRuleSuitepasses.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.