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
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.