Repository navigation
feat: run the operators above a typed Dataset operation natively - #6564
Conversation
Typed Dataset operations (map, flatMap, mapPartitions, mapGroups, cogroup) pass JVM objects between their operators, so they stay on Spark, and today the operators above them stay on Spark too until the next shuffle. Every typed operation ends in SerializeFromObjectExec, whose output is ordinary rows. With the new spark.comet.convert.typedDataset.enabled, CometExecRule puts a CometSparkToColumnarExec above it, so a partial aggregate, a broadcast join or a native shuffle above the operation runs natively. Spark inserts no columnar transitions below a RowToColumnarTransition, so the rule adds them to the subtree under the conversion with Spark's own ApplyColumnarRulesAndInsertTransitions. Without them the typed operation would read its Comet child through Spark's interpreted columnar-to-row path. CometExecRule no longer tags a ColumnarToRowTransition as an unsupported operator, which it did to these transitions on its second pass under AQE. Off by default: it is 1.4-2.2x faster when an aggregate over many groups sits above the typed operation, and slower when the work above is cheap. Adds CometTypedDatasetSuite and CometTypedDatasetBenchmark.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
- Design approach: Add an opt-in
CometSparkToColumnarExecabove supportedSerializeFromObjectExecoutputs. User functions continue running in Spark. - Correctness / compatibility analysis: Spark’s transition rules support the approach across the supported versions. However, I reproduced a join returning 9 rows instead of 100 when conversion creates mixed native/JVM wide-decimal shuffles. The new benchmark also fails strict Spark 3.5 compilation.
- Key design decisions: Reusing Spark’s transition insertion and Comet’s existing Arrow reader keeps the implementation small. Disabling conversion by default appropriately reflects the documented conversion overhead for inexpensive downstream work.
- Implementation sketch: Adds the configuration and conversion helper, registers the new suite in both PR workflows, and supplies tests, a benchmark, and documentation.
- Behavioral changes worth calling out: Compared with latest release branch
branch-1.1at366b157a7583f646ce2d6ca0f87595ebe807e11f, downstream native execution is an intended opt-in change. Default typed Dataset execution remains unchanged. The documented decimal partitioning difference becomes observable row loss in the mixed-path join below. - Suggested improvements: Prevent incompatible mixed decimal partitioning and declare the benchmark row count as
Long. Both findings should be addressed before merge.
Reviewed all nine changed files against base 9c0fde09c075212c99e295c39f44ae67e9ba7414, at full head 2b8706351d514c4189c660027d90d60ba2f859ad. Routed skills: review-comet-pr, review-comet-ffi-pr, review-comet-shuffle-pr, and review-comet-expression-pr. The snapshot and live discussion checks contained no existing reviews, comments, or threads.
Exact-head CI: 21 checks passed, one failed, nine remained in progress, and 42 were skipped. The failure is strict Scala compilation on Spark 3.5. Comet execution suites, Rust tests, benchmark validation, and the requested Spark 4.1 SQL run were still pending.
Local validation: all 10 CometTypedDatasetSuite tests and four additional edge tests passed on Spark 4.1.3/JDK 17. An eight-configuration decimal-join probe reproduced the regression and verified that disabling conversion or forcing JVM shuffle restores the result. JVM sources were built from this head using the base CI native artifact, whose native sources are unchanged by this PR. Other Spark versions were checked against source but not executed locally. Full SQL suites and performance benchmarks were not rerun. Disposable test source was removed and the checkout is clean.
Recommendation: request changes.
| } else { | ||
| val withTransitions = | ||
| ApplyColumnarRulesAndInsertTransitions(Seq.empty, outputsColumnar = false).apply(op) | ||
| convertToComet(withTransitions, CometSparkToColumnarExec).getOrElse(withTransitions) |
There was a problem hiding this comment.
[P1] Prevent conversion from creating incompatible wide-decimal shuffles. With this feature enabled, a supported typed input switches to native shuffle while another typed input containing a retained array<int> column stays on JVM shuffle. Joining their decimal(38,18) keys with AQE disabled then silently loses matching rows: my reproduction returns 9 instead of 100. Conversion-disabled Comet returns all 100. This newly exposes the existing decimal hash difference as incorrect query results. Keep affected exchanges on JVM shuffle until native wide-decimal hashing matches Spark, and cover this mixed-path join.
Evidence: Reproduced on Spark 4.1.3/JDK 17 in CometTestBase with spark.sql.adaptive.enabled=false, spark.sql.autoBroadcastJoinThreshold=-1, spark.sql.shuffle.partitions=10, and spark.comet.shuffle.mode=auto. Define case class L(k: java.math.BigDecimal, v: Long) and case class R(k: java.math.BigDecimal, xs: Seq[Int]). Build l = spark.range(0,100,1,2).map(i => L(new java.math.BigDecimal(i), i)).alias("l") and r = spark.range(0,100,1,2).map(i => R(new java.math.BigDecimal(i), Seq(i.toInt))).alias("r"). Collect l.join(r, col("l.k") === col("r.k")).select(col("l.v"), col("r.xs")). Spark and conversion-disabled Comet return 100 rows. Conversion-enabled Comet returns 9. The executed plan contains left CometNativeShuffle and right CometColumnarShuffle. Setting shuffle mode to jvm restores 100 rows. Spark hashes wide decimals using unscaledValue().toByteArray; native hash_array_decimal! uses fixed-width to_le_bytes().
There was a problem hiding this comment.
Thanks for the repro. It reproduced exactly here: 9 rows instead of 100, with CometNativeShuffle on the left and CometColumnarShuffle on the right.
Fixed in 1ec0e66. A hash shuffle whose stage starts at a typed Dataset conversion now stays on Comet's columnar shuffle when a key is or contains a decimal wider than 18 digits and the shuffle has more than one partition, so the conversion no longer moves it off the path the other side of the join uses. That's the rule #6005 applies to every native shuffle, limited to the shuffles this feature moves, so it becomes redundant once #6005 lands. Your join is now a test (a join on wide decimal keys with an input that is not converted), which fails without the guard. The rule is also listed in native_shuffle.md.
| */ | ||
| object CometTypedDatasetBenchmark extends CometBenchmarkBase { | ||
|
|
||
| private val numRows = 4 * 1024 * 1024 |
There was a problem hiding this comment.
[P2] Declare numRows as a Long so the required strict Spark 3.5 compilation succeeds. It is currently inferred as Int, but both spark.range(numRows) and new Benchmark(name, numRows, ...) require Long. The strict profile promotes those implicit numeric-widening warnings to errors, causing this PR’s CI build to fail. Using 4L * 1024 * 1024 fixes both call sites.
Evidence: Exact-head CI job https://github.com/apache/datafusion-comet/actions/runs/37130634250/job/111225450414 ran ./mvnw -B test-compile -Pspark-3.5 -Pstrict-warnings -DskipTests and failed with CometTypedDatasetBenchmark.scala:161: implicit numeric widening and the same error at line 178. The log ends with two errors found and a failed scala-maven-plugin:testCompile goal.
There was a problem hiding this comment.
Fixed in 1ec0e66. numRows is now 4L * 1024 * 1024, and ./mvnw -B test-compile -Pspark-3.5 -Pstrict-warnings -DskipTests passes locally.
…VM shuffle A typed Dataset conversion moves the shuffle above it from Comet's columnar shuffle, which partitions with Spark's hash, to native shuffle. Native shuffle hashes decimals wider than 18 digits differently from Spark (apache#5994), so a join on such keys with an input that stays on columnar shuffle, for example one with an array<int> column the conversion declines, put matching keys in different partitions and returned 9 rows instead of 100. Such a shuffle now stays on columnar shuffle unless it has one partition, the same rule apache#6005 proposes for every native shuffle. Also declare the benchmark's row count as a Long, which the strict Spark 3.5 compile requires.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
- Design approach: Optionally insert
CometSparkToColumnarExecabove supportedSerializeFromObjectExecoutputs while keeping user functions in Spark. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Checked serializer and transition semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Both earlier findings are addressed: the decimal join returns all expected rows, and exact-head strict Spark 3.5 compilation passes.
- Key design decisions: Reusing Spark’s transition insertion and Comet’s Arrow reader keeps the implementation small. The feature remains off by default, appropriately accounting for conversion overhead when downstream work is inexpensive.
- Implementation sketch: Adds configuration, typed-output conversion, a wide-decimal shuffle guard, regression tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared with
branch-1.1at4f5cf2db1e80bf94db8a4b79a870738d621f1262, downstream native execution is an intended opt-in change. Affected wide-decimal exchanges retain Spark-compatible hashing. The documented exception wrapping applies when conversion is enabled. - Suggested improvements: None at P1/P2 priority.
Reviewed all 11 changed files and both commits against base 9c0fde09c075212c99e295c39f44ae67e9ba7414, at full head 1ec0e66265af56a3b2cfc4e04e6a160efac3c4fc. Confirmed the PR is not a draft and read existing reviews, comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr.
Exact-head CI at 2026-10-03 15:48 UTC: 11 checks passed, 11 were running and 12 were skipped. No failures were reported. Strict Spark 3.5 compilation passed. The Spark SQL 4.1 run, benchmark check, native build and other checks remained unfinished.
Local validation: 105 tests passed on Spark 4.1.3/JDK 17, covering the new suite, planner rules, transition elimination, shuffle fallback and additional edge cases. The decimal-join probe returned 100 rows in all 12 configurations, including AQE with partition coalescing disabled. Built JVM sources from this head using the verified base CI native library; native sources are unchanged. Other Spark versions, full SQL suites and performance benchmarks were not executed locally. Disposable tests were removed and the checkout is clean.
…ests - Use Spark's existsRecursively(DecimalType.isByteArrayDecimalType) for the wide-decimal check, fold the duplicated conversion cases in readsTypedDatasetConversion, and check the config and earlier native shuffle reasons before walking the plan. Add a TODO to drop the guard once apache#5994 lands. - Tests: assert the join's exchanges with checkCometExchange, drop conditions that cannot change the transition assertion, write the Parquet table once for both AQE settings, and take checkConverted's query by value. - Benchmark: import spark.implicits instead of declaring encoders.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
- Design approach: Optionally insert
CometSparkToColumnarExecabove supportedSerializeFromObjectExecoutputs while retaining Spark’s execution of user functions. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Checked serializer, transition, query-stage and decimal-type semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Both earlier findings remain fixed: the mixed-shuffle decimal join preserves its rows, and strict Spark 3.5 compilation passes.
- Key design decisions: Reusing Spark’s transition insertion and the existing Arrow reader keeps the implementation small. The feature remains disabled by default, reflecting the documented conversion overhead for inexpensive downstream work. The decimal guard preserves Spark-compatible hashing for affected exchanges.
- Implementation sketch: Adds configuration, typed-output conversion, shuffle selection safeguards, tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared the affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Native execution above typed operations is an intended opt-in change. Enabling it also exposes the documented existing Arrow-stream exception wrapping. This PR preserves the default conversion policy. - Suggested improvements: None at P1/P2 priority.
Reviewed all 11 changed files and all three commits against base 9c0fde09c075212c99e295c39f44ae67e9ba7414, at full head 61e08895b40d9e33d1ddf54170afb7bf991f7071. Confirmed the PR is not a draft and read existing reviews, comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr.
Exact-head CI: 36 checks passed, 14 were skipped, and none failed or remained running. Passing checks include strict Spark 3.5 compilation, benchmark compilation/lint, Comet execution and shuffle tests, and all requested Spark 4.1 SQL shards. Other Spark SQL versions and macOS runtime checks were skipped.
Local validation: 105 tests passed on Spark 4.1.3/JDK 17, covering typed Datasets, planner rules, transition cleanup, shuffle fallback and additional encoder cases. The decimal-join probe returned all 100 expected rows in all 12 configurations, including AQE with partition coalescing disabled. Built current JVM sources using the verified exact-head native CI artifact. The wrapper’s read-only default cache required using cached Maven directly. Other Spark versions, full SQL suites and performance benchmarks were not executed locally. Disposable test source was removed and the checkout is clean.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
- Design approach: Optionally convert
SerializeFromObjectExecoutput to Arrow while retaining Spark execution of user functions. - Correctness / compatibility analysis: Reproduced one new P2 issue: batching evaluates user code beyond a downstream
limit, causing a query that succeeds in Spark to fail. The earlier decimal-join and strict-compilation findings are fixed. Checked serializer, transition and limit semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. - Key design decisions: Reusing Spark’s transition insertion and Comet’s Arrow reader keeps the implementation small. The public configuration remains disabled by default, keeping the documented conversion overhead opt-in.
- Implementation sketch: Adds configuration, typed-output conversion, transition handling, a wide-decimal shuffle guard, tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Native execution above typed operations is an intended opt-in change. Failing a limited query because of later rows is an unintended compatibility regression. - Suggested improvements: Preserve row-level short-circuiting before batching, or retain Spark execution for affected limit pipelines. Add the reproduced regression case.
Reviewed the full 11-file diff from base ff468b197f53369d6f1086d691e6d34fe32e4caf to head afb7ebfda74a227d7f421b4f5227ec98fc71764e. Confirmed the PR is not a draft and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr.
Exact-head CI at 2026-10-04 17:10 UTC: 18 checks passed, 12 were skipped and eight remained running. No failures were reported. Strict Spark 3.5 compilation and benchmark compilation/lint passed. Comet runtime suites, Rust tests, TPC result checks and the requested Spark 4.1 SQL run remained unfinished.
Local validation: Ran 17 focused tests on Spark 4.1.3/JDK 17. Sixteen passed, including all 11 PR tests. The disposable limit regression test failed only with conversion enabled, with AQE both on and off. The decimal-join probe returned all 100 rows in all 12 configurations. Used the verified base CI native artifact, whose native sources match this head. Other Spark versions were source-checked only. Full SQL suites and performance benchmarks were not run locally. Disposable test source was removed and the checkout is clean.
Recommendation: request changes for the new P2 finding.
| } else { | ||
| val withTransitions = | ||
| ApplyColumnarRulesAndInsertTransitions(Seq.empty, outputsColumnar = false).apply(op) | ||
| convertToComet(withTransitions, CometSparkToColumnarExec).getOrElse(withTransitions) |
There was a problem hiding this comment.
[P2] Preserve short-circuiting when a limit consumes typed output. With spark.comet.convert.typedDataset.enabled=true, map(...).limit(1) fills an Arrow batch before returning its first row. A user function that throws on row 30 therefore fails the query, although Spark and conversion-disabled Comet return the first row successfully. The resulting plan is CometCollectLimit -> CometSparkRowToColumnar -> SerializeFromObject. Please preserve row-level limiting before batching where valid, or decline conversion for this pipeline, and add a regression test.
Evidence: Reproduced at the reviewed head in CometTestBase on Spark 4.1.3/JDK 17, with AQE both false and true: spark.range(0, 100, 1, 1).map { i => if (i == 30L) throw new IllegalArgumentException("unexpected evaluation of row 30"); i + 1L }.toDF().limit(1).collect(). Spark and Comet with typed conversion disabled return [Row(1)]. Enabling conversion throws SparkException caused by that IllegalArgumentException, with RowArrowReader.loadNextBatch in the stack. Spark’s CollectLimitExec.executeCollect calls child.executeTake(limit), whereas the inserted Arrow reader consumes a batch before the limit can stop it. The six-configuration probe failed only in the two conversion-enabled cases.
There was a problem hiding this comment.
Fixed in 2737816, and generalized in c202f51. The conversion leaves a typed operation's output unconverted when a limit above it can stop reading early, unless an operator that reads all of its input, such as an exchange, a sort or a hash aggregate, sits in between. Your query is the test a limit does not evaluate typed Dataset rows beyond the result, with AQE on and off.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
- Design approach: Optionally convert
SerializeFromObjectExecoutput to Arrow while retaining Spark execution of user functions. - Correctness / compatibility analysis: No additional introduced P1/P2 issues found within this review. The existing P2 limit regression remains reproducible:
map(...).limit(1)evaluates a later throwing row when conversion is enabled, whereas Spark and conversion-disabled Comet return the first row. Earlier decimal-join and strict-compilation findings are addressed. Checked serializer, transition, limit and decimal-hashing semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. - Key design decisions: Reusing Spark’s transition insertion and Comet’s Arrow reader keeps the implementation small. The public configuration remains disabled by default, appropriately making the documented conversion overhead opt-in.
- Implementation sketch: Adds configuration, typed-output conversion, transition handling, a wide-decimal shuffle guard, tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Native execution above typed operations is an intended opt-in change. Evaluating throwing rows beyond a downstream limit is an unintended compatibility regression. - Suggested improvements: Resolve the existing limit finding by preserving row-level short-circuiting before batching or declining conversion for affected pipelines, with regression coverage.
Reviewed all 11 changed files in the full diff from base ff468b197f53369d6f1086d691e6d34fe32e4caf to head afb7ebfda74a227d7f421b4f5227ec98fc71764e. Confirmed the PR is not a draft and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr.
Exact-head CI at 2026-10-04 17:18 UTC: 23 checks passed, 12 were skipped and 13 remained running. No failures were reported. Strict Spark 3.5 compilation, benchmark compilation/lint, shuffle tests and TPC-H checks passed. Execution, expression, scan, TPC-DS and Spark 4.1 SQL checks remained pending.
Local validation: Rebuilt JVM sources and ran 17 focused tests on Spark 4.1.3/JDK 17 using the exact-head native CI artifact. Sixteen passed, including all 11 PR tests. The existing limit reproduction failed only with conversion enabled, with AQE both on and off. The decimal-join probe returned all 100 rows in all 12 configurations. Other Spark versions were source-checked only. Full SQL suites and performance benchmarks were not run locally. Disposable test source was removed and the checkout is clean.
Recommendation: request changes for the unresolved existing P2. No duplicate finding is added.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
- Design approach: Optionally convert
SerializeFromObjectExecoutput to Arrow while retaining Spark execution of user functions. - Correctness / compatibility analysis: The earlier decimal-join, physical-limit and strict-compilation findings are addressed. One additional P2 is reproducible:
mapPartitions(_.take(1))evaluates throwing upstream rows when conversion is enabled. Checked serializer, iterator, limit and transition semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. - Key design decisions: Reusing Spark’s transition insertion and Comet’s existing Arrow reader keeps the implementation small. Disabling the feature by default appropriately makes the documented conversion overhead opt-in.
- Implementation sketch: Adds configuration, typed-output conversion, transition handling, decimal-shuffle and physical-limit guards, tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Native execution above typed operations is an intended opt-in change. Evaluating rows that a downstream iterator never requests is an unintended compatibility regression. - Suggested improvements: Preserve lazy consumption across typed iterator boundaries and add coverage for the reproduced
mapPartitionscase.
Reviewed all 11 changed files in the full diff from base ff468b197f53369d6f1086d691e6d34fe32e4caf to head 2737816792153ea2bf2161f48b699dd253a55087. Confirmed the PR is not a draft and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr.
Exact-head CI at 2026-10-04 18:09 UTC: 21 checks passed, 12 were skipped and seven remained running. No failures were reported. Strict Spark 3.5 compilation, benchmark compilation/lint, native compilation and Rust tests passed. Comet runtime suites, TPC result checks and the requested Spark 4.1 SQL build remained unfinished.
Local validation: Ran 19 focused tests on Spark 4.1.3/JDK 17. Eighteen passed, including all 12 PR tests. The iterator regression failed only with conversion enabled, with AQE both on and off. The decimal join returned all 100 rows in all 12 configurations, and the original limit reproduction passed all six configurations. Used a checksum-verified prior-head CI native library whose sources are unchanged at this head. Other Spark versions were source-checked only. Full SQL suites and performance benchmarks were not run locally. Disposable test source was removed and the checkout is clean.
Recommendation: request changes for the new P2 finding.
| "Comet does not convert the output of a typed Dataset operation below a limit " + | ||
| "because Arrow batching could evaluate rows beyond Spark's row-level limit") | ||
| } else { | ||
| convertTypedDatasetOutput(op) |
There was a problem hiding this comment.
[P2] Preserve lazy consumption before downstream MapPartitionsExec consumers. The physical-limit guard does not cover iterator operations such as mapPartitions(_.take(1)). With typed conversion enabled, an intervening native filter batches the earlier typed output and evaluates row 30 before returning the first result. A user function throwing on row 30 therefore fails a query that succeeds in Spark and conversion-disabled Comet. Please decline conversion across these lazy iterator boundaries unless full consumption is guaranteed, or otherwise preserve row-level consumption, and add this regression case.
Evidence: Reproduced at this head on Spark 4.1.3/JDK 17: spark.range(0, 100, 1, 1).map { i => if (i == 30L) throw new IllegalArgumentException("unexpected evaluation of row 30"); i + 1L }.filter(col("value") > 0L).mapPartitions(_.take(1)).toDF().collect(). With AQE both false and true, Spark and Comet with spark.comet.convert.typedDataset.enabled=false return [Row(1)]. Setting it to true throws SparkException caused by CometNativeException wrapping that IllegalArgumentException. The failing plan contains MapPartitions -> DeserializeToObject -> CometColumnarToRow -> CometFilter -> CometSparkRowToColumnar -> SerializeFromObject, with no physical limit node. Spark’s MapPartitionsExec passes a lazy mapped iterator to the user function, whereas RowArrowReader.loadNextBatch consumes upstream rows to fill a batch. The six-configuration probe failed only in the two conversion-enabled cases.
There was a problem hiding this comment.
Fixed in c202f51. A MapPartitionsExec above the typed operation now counts as a reader that can stop early, as a limit does. Your query is the test a mapPartitions function does not evaluate typed Dataset rows it never reads, with AQE on and off.
…ps early The limit guard covered only physical limits, so a mapPartitions function such as `_.take(1)`, or code reading Dataset.rdd, still read typed rows through an Arrow batch and ran the user function on rows that Spark never reaches. Walk down from each reader that can stop early (a limit, a top-k over sorted input, a MapPartitionsExec, and a DeserializeToObjectExec at the plan root) and stop at an operator that reads all of its input first, so a limit above an aggregate keeps the conversion.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
- Design approach: Optionally convert
SerializeFromObjectExecoutput to Arrow while retaining Spark execution of user functions. - Correctness / compatibility analysis: The earlier decimal-join, strict-compilation, physical-limit and
mapPartitionsfindings are addressed. One new P2 remains: an optimized.rdd.take(1)pipeline bypasses the new safeguard and evaluates throwing rows Spark never requests. Checked relevant serializer, iterator, optimizer and transition semantics across Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. - Key design decisions: Reusing Spark’s transition rules and Comet’s existing Arrow reader limits duplication. Keeping conversion disabled by default appropriately accounts for the documented overhead with inexpensive downstream operations.
- Implementation sketch: Adds configuration, conversion and transition handling, partial-consumption and decimal-shuffle safeguards, tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Native execution above typed operations is an intended opt-in change. The reproduced RDD failure is an unintended compatibility regression. Arrow-stream exception wrapping is documented. - Suggested improvements: Recognize optimized RDD roots when protecting lazy consumption and cover the reproduced second-
mapcase.
Reviewed the full 11-file diff from base ff468b197f53369d6f1086d691e6d34fe32e4caf to head c202f515888da0a6a81f526e4abe152ed4c724c7. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr.
Exact-head CI at 2026-10-04 19:14 UTC: 24 checks passed, 12 were skipped and 11 remained running. No failures were reported. Strict Spark 3.5 compilation, benchmark checks, Rust tests, shuffle, scans and TPC result checks passed. Execution, expression and nine Spark 4.1 SQL jobs remained unfinished.
Local validation: Ran 133 tests on Spark 4.1.3/JDK 17. Of these, 132 passed, including all 15 PR tests. The new RDD regression failed only with conversion enabled, with AQE both on and off. Earlier decimal, limit and iterator reproductions passed. Built current JVM sources using the checksum-verified exact-head native CI library. Other Spark versions were source-checked only. Full SQL suites and performance benchmarks were not run locally. Disposable tests were removed and the checkout is clean.
Recommendation: request changes for the new P2 finding.
| // `Dataset.rdd` plans a `DeserializeToObjectExec` at the root, and the RDD's own code | ||
| // decides how much of it to read, as `take(1)` does. | ||
| val rootReader = plan match { | ||
| case _: DeserializeToObjectExec => Some("code reading Dataset.rdd") |
There was a problem hiding this comment.
[P2] Preserve RDD laziness after Spark eliminates serialization. Dataset.rdd does not always retain a DeserializeToObjectExec root: when the Dataset ends with another map, Spark’s EliminateSerialization removes that deserializer and leaves MapElementsExec. Consequently, map(f).filter(...).map(g).rdd.take(1) bypasses this guard and converts the earlier serializer. If f throws on row 30, conversion-enabled Comet fails although Spark and conversion-disabled Comet return the first result successfully. Could the safeguard recognize these optimized RDD roots and preserve row-level consumption? Add this case alongside the current RDD regression test.
Evidence: Reproduced at this head on Spark 4.1.3/JDK 17: spark.range(0, 100, 1, 1).map { i => if (i == 30L) throw new IllegalArgumentException("unexpected evaluation of row 30"); i + 1L }.filter(col("value") > 0L).map(i => i + 1L).rdd.take(1). With AQE both false and true, Spark and Comet with spark.comet.convert.typedDataset.enabled=false return Array(2L). Enabling conversion throws SparkException caused by CometNativeException wrapping that IllegalArgumentException, with RowArrowReader.loadNextBatch in the stack. The RDD execution plan is rooted at MapElements, above DeserializeToObject -> CometColumnarToRow -> CometFilter -> CometSparkRowToColumnar -> SerializeFromObject. Spark’s EliminateSerialization rule removes the outer deserialize/serialize pair across all five supported versions. The six-configuration probe failed only in the two conversion-enabled cases.
There was a problem hiding this comment.
Confirmed: with the conversion on, map(f).filter(...).map(g).rdd.take(1) failed on row 30, with MapElements at the root of the plan. 1e23881 takes any root that produces objects as code reading Dataset.rdd: an ObjectProducerExec, or a filter or project over one. That covers the DeserializeToObjectExec that Dataset.rdd adds, the typed operation left at the root once EliminateSerialization drops it, and a typed filter Spark pulls above that operation. A Dataset's own plan ends in rows, so no other plan has such a root.
The RDD test now runs your query, and the same query with a typed filter at the end, with AQE on and off. Your query fails with the old check, and the typed filter one fails without the walk through the filter.
|
@copilot resolve the merge conflicts on this branch. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
- Design approach: Optionally convert
SerializeFromObjectExecoutput to Arrow while retaining Spark execution of user functions. - Correctness / compatibility analysis: No additional introduced P1/P2 issues found within this review. The existing P2 RDD-laziness finding remains reproducible. With AQE both on and off,
map(f).filter(...).map(g).rdd.take(1)returns2in Spark and conversion-disabled Comet, but conversion-enabled Comet evaluates the throwing row 30 and fails. Earlier decimal-join, physical-limit,mapPartitionsand strict-compilation findings are addressed. - Key design decisions: Reusing Spark’s transition insertion and Comet’s Arrow reader limits duplication. Keeping conversion disabled by default appropriately accounts for the documented overhead with inexpensive downstream work.
- Implementation sketch: Adds configuration, typed-output conversion, transition handling, partial-consumption and decimal-shuffle safeguards, tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Native execution above typed operations is an intended opt-in change. The reproduced RDD failure remains an unintended compatibility regression. Default typed-output conversion stays disabled. - Suggested improvements: Resolve the existing RDD finding by recognizing optimized object-producing roots when protecting lazy consumption, and cover the second-
mapreproduction.
Reviewed all 11 changed files in the full diff from base ff468b197f53369d6f1086d691e6d34fe32e4caf to head c202f515888da0a6a81f526e4abe152ed4c724c7. Confirmed non-draft status and read existing reviews, comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr. Checked relevant Spark serializer, iterator, optimizer and transition sources across all five supported versions.
Exact-head CI: 36 checks passed, 14 were skipped, and none failed or remained running. Passing checks include strict Spark 3.5 compilation, benchmark checks, Rust tests, Comet runtime suites and all nine Spark 4.1 SQL shards. Other Spark SQL versions and macOS runtime checks were skipped.
Local validation: Rebuilt JVM sources and ran 23 focused tests on Spark 4.1.3/JDK 17 using the checksum-verified exact-head native CI artifact. Twenty-two passed, including all 15 PR tests. The sole failure confirmed the existing RDD regression. The decimal join returned all 100 rows in all 12 configurations. Other Spark versions were source-checked only. Full SQL suites and performance benchmarks were not rerun locally. Disposable test sources were removed and the checkout is clean.
Recommendation: request changes for the unresolved existing P2. No duplicate finding is added.
…w-to-columnar # Conflicts: # spark/src/main/scala/org/apache/comet/CometConf.scala # spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
- Design approach: Optionally convert supported
SerializeFromObjectExecoutput to Arrow while retaining Spark execution of user functions. - Correctness / compatibility analysis: No additional introduced P1/P2 issues found within this review. The existing P2 RDD-laziness finding remains reproducible. With AQE on and off,
map(f).filter(...).map(g).rdd.take(1)returns2in Spark and conversion-disabled Comet, but conversion-enabled Comet evaluates the throwing row 30 and fails. Earlier decimal-join, physical-limit,mapPartitionsand strict-compilation findings are addressed. Checked relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. - Key design decisions: Reusing Spark’s transition insertion and Comet’s Arrow reader limits duplication. Keeping conversion disabled by default appropriately accounts for the measured overhead with inexpensive downstream work.
- Implementation sketch: Adds configuration, typed-output conversion, transition handling, partial-consumption and decimal-shuffle safeguards, tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Downstream native execution is an intended opt-in change. The reproduced RDD failure remains an unintended compatibility regression. - Suggested improvements: Resolve the existing RDD finding by recognizing optimized object-producing roots when protecting lazy consumption, and cover the second-
mapreproduction.
Reviewed all 11 changed files in the full diff from base 965c8bbe289ff850b614c8c3833ed7802134115b to head bd36c6b57d4b5f5ee3a7b47d05b4004ee0955d29. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr.
Exact-head CI at 2026-10-05 13:51 UTC: 17 checks passed, 12 were skipped and three remained running. No failures were reported. Strict Spark 3.5 compilation and benchmark compilation/lint passed. Native compilation, Rust tests and the Spark 4.1 SQL build remained running. Comet runtime and SQL test verdicts were not yet available.
Local validation: Built this head’s JVM sources and ran 23 focused tests on Spark 4.1.3/JDK 17 using a checksum-verified base CI native library whose sources are unchanged in this PR. Twenty-two passed, including all 15 PR tests. The sole failure confirmed the existing RDD regression. The decimal join returned all 100 rows in all 12 configurations. Other Spark versions were source-checked only. Full SQL suites, Rust tests and performance benchmarks were not run locally. Disposable test sources were removed and the checkout is clean.
Recommendation: request changes for the unresolved existing P2. No duplicate finding is added.
The guard for code reading Dataset.rdd looked for the DeserializeToObjectExec that Dataset.rdd puts at the root of the plan. When the Dataset ends in a typed operation such as map, Spark's EliminateSerialization drops that deserializer together with the operation's serializer, so the root is the operation itself, or a typed filter over it. The output of an earlier typed operation was then still converted, and map(f).filter(...).map(g).rdd.take(1) ran f on rows that Spark never reaches. The guard now takes any root that produces objects, under a filter or a project, as code reading Dataset.rdd. A Dataset's own plan ends in rows, so no other plan has such a root.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
- Design approach: Optionally convert supported
SerializeFromObjectExecoutput to Arrow while retaining Spark execution of user functions. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Earlier decimal-join, strict-compilation, limit, iterator and optimized RDD findings are addressed. Checked relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0.
- Key design decisions: Reusing Spark’s transition insertion and Comet’s Arrow reader limits duplication. Keeping conversion disabled by default accounts for the documented overhead with inexpensive downstream work. Partial-consumption guards preserve laziness, and the decimal-shuffle guard preserves compatible partitioning.
- Implementation sketch: Adds configuration, conversion and transition handling, safeguards, tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Downstream native execution is an intended opt-in change. Default typed-output conversion remains disabled. Enabling conversion exposes the documented existing Arrow-stream exception wrapping. - Suggested improvements: None at P1/P2 priority.
Reviewed all 11 changed files in the full diff from base 965c8bbe289ff850b614c8c3833ed7802134115b to head 1e238811fe3c04d61c940c9d06fcc5ac708925a1. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr and review-comet-expression-pr.
Exact-head CI at 2026-10-05 14:20 UTC: 20 checks passed, 12 were skipped, seven were running and nine were queued. No failures were reported. Strict Spark 3.5 compilation and benchmark compilation/lint passed. Comet runtime suites, Rust tests, TPC result checks and nine Spark 4.1 SQL shards remained unfinished. Other Spark SQL versions and macOS runtime checks were skipped.
Local validation: Built this head’s JVM sources on Spark 4.1.3/JDK 17 using the checksum-verified exact-head native CI artifact. All 130 selected repository tests passed, including the 15 PR tests. Four disposable probes also passed after correcting one probe’s expected Spark behavior. Earlier laziness reproductions passed with AQE on and off, and the decimal join returned all 100 rows in all 12 configurations. Other Spark versions were source-checked only. Full SQL suites, Rust tests and performance benchmarks were not run locally. Local formatting/enforcer checks were skipped. Disposable test source was removed and the checkout is clean.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
- Design approach: Optionally convert supported
SerializeFromObjectExecoutput to Arrow while retaining Spark execution of user functions. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Earlier decimal-join, strict-compilation, limit, iterator and optimized RDD findings are addressed. Checked relevant Spark semantics across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0.
- Key design decisions: Reusing Spark’s transition insertion and Comet’s Arrow reader limits duplication. Disabling conversion by default accounts for its documented overhead with inexpensive downstream work. The partial-consumption and decimal-shuffle guards preserve existing behavior.
- Implementation sketch: Adds configuration, conversion and transition handling, safeguards, regression tests, a benchmark, workflow registration and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Downstream native execution is an intended opt-in change. Default typed-output conversion remains disabled. Enabling conversion exposes the documented existing Arrow-stream exception wrapping. - Suggested improvements: None at P1/P2 priority.
Reviewed the full 11-file diff from base 965c8bbe289ff850b614c8c3833ed7802134115b to head 1e238811fe3c04d61c940c9d06fcc5ac708925a1. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-ffi-pr, review-comet-shuffle-pr and review-comet-expression-pr.
Exact-head CI at 2026-10-05 14:29 UTC: 24 checks passed, 12 were skipped and 12 remained running. No failures were reported. Strict Spark 3.5 compilation, benchmark checks, Rust tests, scan/shuffle suites and TPC-H checks passed. Execution/expression suites, TPC-DS checks and all nine Spark 4.1 SQL shards remained unfinished.
Local validation: All 15 PR tests and six disposable probes passed across validation runs on Spark 4.1.3/JDK 17. Earlier laziness reproductions passed with AQE on and off, and the decimal join returned all 100 rows in 12 configurations. One probe initially assumed sorting evaluated each row once. Comparison with Spark confirmed identical range-sampling counts in all modes. Used current JVM sources and the checksum-verified exact-head native CI artifact. Other Spark versions were source-checked only. Full SQL suites, Rust tests and performance benchmarks were not run locally. Local formatting/enforcer checks were skipped. Disposable tests were removed and the checkout is clean.
comphead
left a comment
There was a problem hiding this comment.
Thanks for the thorough tests and the benchmark. I read the limit, mapPartitions and Dataset.rdd guards and the earlier findings in this thread, and I did not find a gap in them. The main item is the struct column failure in the first inline comment, which I ran. The other comments come from reading the source.
One request without an anchor: the description reports the four tests added after the five-version run (the limit, mapPartitions, Dataset.rdd and limit-above-aggregate cases) on Spark 3.5 and 4.1 only. They depend on Spark's optimizer (EliminateSerialization, limit planning), and the Comet suites run on 3.4, 4.0 and 4.2 only in the nightly. Could you add the run-all-spark-profiles label so they run before this is queued?
| } else { | ||
| val withTransitions = | ||
| ApplyColumnarRulesAndInsertTransitions(Seq.empty, outputsColumnar = false).apply(op) | ||
| convertToComet(withTransitions, CometSparkToColumnarExec).getOrElse(withTransitions) |
There was a problem hiding this comment.
A native shuffle placed directly over this conversion reads it through executeColumnar() and ColumnarBatchArrowReader, which closes each batch. RowArrowReader reuses its vectors, and closing a struct vector drops its children, so a struct column fails from the second batch. #6607 hits the same thing and fixes it by reading a CometNativeArrowSource child as an Arrow stream.
I replayed this with the CometExecRule, CometShuffleExchangeExec and CometConf from this head over a local build of main (Spark 4.1.3, spark.comet.convert.typedDataset.enabled=true, spark.comet.batchSize=16):
case class Rec(a: Int, b: String)
case class Out(id: Int, inner: Rec)
spark.range(0, 200, 1, 2).map(i => Out(i.toInt, Rec(i.toInt, s"s$i"))).repartition(col("id")).collect()This fails with CometNativeException: ... no more field nodes for field a, with AQE on and off. orderBy("id"), the left side of a shuffled join, and a union that feeds a shuffle fail the same way. The plan is CometExchange ... CometNativeShuffle over CometSparkRowToColumnar over SerializeFromObject. The same repartition returns all 200 rows with one batch, with flat columns, with an array<string> column, or with the conversion off.
No test here puts a native shuffle directly over the conversion. The struct test reads only id downstream, so Spark's ObjectSerializerPruning removes inner and tags from the serializer before the conversion sees them (a plain explain of that query shows SerializeFromObject [... AS id]). Could you add a struct case with a shuffle directly above the conversion and a small spark.comet.batchSize? It will fail until #6607 lands, so please either land that first or decline struct columns in convertTypedDatasetOutput until then.
| // A top-k reads only its first rows when its input is already sorted. | ||
| case topK: TakeOrderedAndProjectExec | ||
| if SortOrder.orderingSatisfies(topK.child.outputOrdering, topK.sortOrder) => | ||
| Some("a limit") |
There was a problem hiding this comment.
Can this arm ever tag a SerializeFromObjectExec? Typed operators report no outputOrdering (SerializeFromObjectExec does not override it), so an ordering that satisfies the top-k comes from a SortExec, and the SortExec arm below resets the search before it reaches the serializer. I could not build a plan where this arm changes a tag, and no test reaches it. If there is none, dropping the arm, its comment and the SortOrder import leaves TakeOrderedAndProjectExec as a plain barrier. If there is one, a test would pin it.
| // wider than 18 digits differently from Spark, so a join with an input that is still on | ||
| // the columnar shuffle would put matching keys in different partitions. Leave such a | ||
| // shuffle where it was. A single partition hashes nothing. | ||
| // TODO: remove once native hashing matches Spark for wide decimals (#5994). |
There was a problem hiding this comment.
#6005 rejects native hashing of wide decimal keys for every native shuffle, and #6607 carries its own copy of this guard (hashesDifferentlyFromSpark). Once #6005 lands, this block, readsTypedDatasetConversion, the native_shuffle.md bullet and the join test are no longer needed. This TODO names only #5994, so the guard would stay behind after #6005. Could it name #6005 as well, so whoever lands that removes this guard in the same change? #6005 also edits the same section of native_shuffle.md, so one of the two will need a rebase.
|
|
||
| import spark.implicits._ | ||
|
|
||
| private case class Arm(name: String, confs: Seq[(String, String)], check: SparkPlan => Unit) |
There was a problem hiding this comment.
Arm, runArm and the check-then-time loop in runCometBenchmark repeat CometRangeBenchmark, where the Arm signature is identical. #6607 adds a third copy in CometShuffleInputConversionBenchmark. Could Arm and one runArms(name, numRows, arms, query) helper move into CometBenchmarkBase in whichever of the two PRs lands first, so this file keeps only its cases and plan checks?
| The same types apply to the output of typed `Dataset` operations, such as `map`, which Comet | ||
| converts when `spark.comet.convert.typedDataset.enabled=true`. A column of any other type keeps | ||
| the operators above the typed operation on Spark. |
There was a problem hiding this comment.
The "Other Spark inputs" list above names each spark.comet.convert.* conversion of a Spark operator's output, and this one is missing from it. Could it get a bullet there, for example spark.comet.convert.typedDataset.enabled: the output of typed Dataset operations such as map and mapPartitions? The sentence about column types can stay here. Someone looking for what Comet can convert will not find the new config otherwise.
|
I reproduced the struct failure, but it isn't new in this PR, so I filed #6685 for it. Main fails the same way today with On the profiles, I ran |
Which issue does this PR close?
Closes #5710. Part of #5572. Supersedes #5714.
Rationale for this change
A typed
Datasetoperation such asds.map(f)plans aDeserializeToObject/MapElements/SerializeFromObjectisland that has to run in Spark, because it passes JVM objects between its operators. Comet only takes over again at the next shuffle, so whatever sits between the island and that shuffle stays on Spark too. Fords.map(f).groupBy(k).agg(...)that is the partial aggregate, which is expensive when there are many groups. A broadcast join on the probe side is in the same position.#5714 tried to fuse the island into one projection in the JVM codegen dispatcher. Since #6542 the dispatcher only calls into Spark's own classes, so that approach is blocked. When I measured the alternative in this PR against it, the conversion got 93-100% of the fused speedup in the cases where either one pays off. It is also much simpler, and it covers every typed operation rather than only
map.What changes are included in this PR?
Every typed operation ends in
SerializeFromObjectExec, whose output is ordinary rows. With the newspark.comet.convert.typedDataset.enabled,CometExecRuleputs aCometSparkToColumnarExecabove it, so the operators above the typed operation can run natively. The operation itself, including the user function, runs in Spark exactly as before. This coversmap,flatMap,mapPartitions,groupByKey(...).mapGroups,cogroup, and a Dataset built from an RDD of objects. If nothing native consumes the output,EliminateRedundantTransitionsremoves the conversion again, so a typed operation at the top of a plan is unchanged.Spark inserts no columnar transitions below a
RowToColumnarTransition, andCometSparkToColumnarExecis one. Above a leaf that does not matter. Here the typed operation's own operators sit below the conversion, and without transitions they read their Comet child throughCometExec.doExecute, which is Spark's interpreted columnar-to-row path. That gives the right answer slowly, and in the benchmark it made whole queries 1.5-2.4x slower. The rule therefore applies Spark's ownApplyColumnarRulesAndInsertTransitionsto the subtree, asCometRule.buildPreviewalready does, andEliminateRedundantTransitionsswaps in Comet's columnar-to-row as usual. Spark's rule leaves existing transitions alone, which matters becauseCometExecRuleruns over the same plan twice under AQE.That second pass used to tag the inserted
ColumnarToRowExecwith "ColumnarToRow is not supported", soCometExecRuleno longer reports aColumnarToRowTransitionas an operator it failed to convert. A column type that Spark-to-Comet conversion does not support, such asarray<int>, keeps the operators above the typed operation on Spark and records a fallback reason that names the column. Spark computes a typed operation's rows one at a time, as they are read, while the conversion fills a whole Arrow batch first. So where something above can stop reading early, namely a limit, a top-k over sorted input, amapPartitionsfunction, or code readingDataset.rdd, the conversion declines with a fallback reason, unless an operator that reads all of its input first, such as an exchange, a sort or a hash aggregate, sits in between. Otherwise the user function would run on rows that Spark never reaches. Code readingDataset.rddis recognized by a root that produces objects, since Spark'sEliminateSerializationcan drop the deserializerDataset.rddadds.It is off by default because whether it pays depends on the work above the typed operation, which the planner cannot see. Here is
CometTypedDatasetBenchmarkonkube2, an AMD Ryzen 9 7950X, with 4Mi rows,local[1], AQE off and one shuffle partition. Times are best of the iterations in ms, and the last column repeats the default Comet arm as a noise check.With an aggregate over many groups above the typed operation, it is up to 1.3x faster than today's default on this host. With a cheap aggregate it is slower, 0.5-0.7x, because Spark compiles the typed operation, the filter and a small aggregate into one loop, while the conversion writes every row to Arrow first.
Two behaviors worth knowing when it is on:
CometNativeException: C Data interface error: java.lang.IllegalArgumentException: ...rather than as the original exception. Every Spark-to-Arrow input wraps errors this way today (Spark errors after the first batch of a JVM input stream reach the user as CometNativeException #6234), and the open fix: rethrow the exception a JVM input throws instead of a CometNativeException #6243 fixes it for all of them.The user guide's operator page and its section on Spark-to-Comet conversion types describe the new config. The contributor guide's paragraph on typed operators suggested the dispatcher fuse, so it now describes the conversion and why the fuse was dropped.
How are these changes tested?
New
CometTypedDatasetSuitehas 15 tests, registered in both PR workflows. Most of them check the answer against Spark, that every operator above the typed operation is native, and that the conversion sits directly onSerializeFromObjectExec. They cover:mapfollowed by an aggregate with AQE on and off, andflatMap,mapPartitions,mapGroups,cogroup, and an RDD of objects.CometBroadcastHashJoinExec.decimal(38,18)keys between a converted input and one that cannot convert, which keeps all its rows because both shuffles stay columnar.Option, nested struct andarray<string>fields, and anarray<int>column that declines with its fallback reason.mapPartitionsfunction and code readingDataset.rdd, including through a Dataset that ends in a typed operation, never running the user function on rows they don't read, with AQE on and off. A limit above an aggregate keeps the conversion.Every converted plan is also checked for row operators reading a columnar child without a transition, and for spurious
ColumnarToRowfallback reasons. Dropping the transition insertion fails 5 tests, and dropping theColumnarToRowTransitionchange fails 4.The original 11-test suite passes locally on Spark 3.4, 3.5, 4.0, 4.1 and 4.2. The full 15-test suite passes locally on Spark 3.5, with strict warnings, and 4.1. On 4.1 these also pass:
CometExecSuite,CometExecRuleSuite,EliminateRedundantTransitionsSuite,RevertNativeForTransitionHeavyStagesSuite,CometInMemoryCacheSuite,CometRangeExecSuite,CometJoinSuite, and both TPC-DS plan stability suites with no approved plan changes. Semantic scalafix (3.4), scalastyle, spotless and prettier are clean.I also ran Spark's own typed
Datasetsuites from the 4.1.3 tests jar, with Comet enabled through system properties:DatasetSuite,DatasetPrimitiveSuite,DatasetAggregatorSuite,DatasetOptimizationSuite,DatasetCacheSuiteandDatasetSerializerRegistratorSuite, 281 tests. With the conversion off and on, the same 278 pass and the same 3 fail: the two TIME-type tests, which need Spark's test-only confs, andgroupBy.as, which asserts on Spark's own plan nodes. A query listener showed that 22 of the suites' queries ran with the conversion when it was on.