Skip to content

feat: run the operators above a typed Dataset operation natively - #6564

Merged
andygrove merged 8 commits into
apache:mainfrom
andygrove:feat/typed-dataset-row-to-columnar
Oct 5, 2026
Merged

andygrove merged 8 commits into
apache:mainfrom
andygrove:feat/typed-dataset-row-to-columnar

Conversation

@andygrove

@andygrove andygrove commented Oct 3, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #5710. Part of #5572. Supersedes #5714.

Rationale for this change

A typed Dataset operation such as ds.map(f) plans a DeserializeToObject / MapElements / SerializeFromObject island 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. For ds.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 new spark.comet.convert.typedDataset.enabled, CometExecRule puts a CometSparkToColumnarExec above 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 covers map, flatMap, mapPartitions, groupByKey(...).mapGroups, cogroup, and a Dataset built from an RDD of objects. If nothing native consumes the output, EliminateRedundantTransitions removes the conversion again, so a typed operation at the top of a plan is unchanged.

Spark inserts no columnar transitions below a RowToColumnarTransition, and CometSparkToColumnarExec is 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 through CometExec.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 own ApplyColumnarRulesAndInsertTransitions to the subtree, as CometRule.buildPreview already does, and EliminateRedundantTransitions swaps in Comet's columnar-to-row as usual. Spark's rule leaves existing transitions alone, which matters because CometExecRule runs over the same plan twice under AQE.

That second pass used to tag the inserted ColumnarToRowExec with "ColumnarToRow is not supported", so CometExecRule no longer reports a ColumnarToRowTransition as an operator it failed to convert. A column type that Spark-to-Comet conversion does not support, such as array<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, a mapPartitions function, or code reading Dataset.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 reading Dataset.rdd is recognized by a root that produces objects, since Spark's EliminateSerialization can drop the deserializer Dataset.rdd adds.

It is off by default because whether it pays depends on the work above the typed operation, which the planner cannot see. Here is CometTypedDatasetBenchmark on kube2, 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.

case Spark Comet (default) Comet, converted default repeat
map -> group by 100 keys 269 229 346 222
map -> group by 1M keys 1060 927 706 940
map -> filter -> group by 100 keys 211 167 355 166
map -> filter -> group by 1M keys 760 654 621 655
map -> group by long key, 1M keys 798 663 533 654
map, 4 columns -> group by 1M keys 1253 1123 926 1135
mapPartitions -> group by 1M keys 1194 1103 848 1125

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:

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 CometTypedDatasetSuite has 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 on SerializeFromObjectExec. They cover:

  • map followed by an aggregate with AQE on and off, and flatMap, mapPartitions, mapGroups, cogroup, and an RDD of objects.
  • A broadcast join above the typed operation running as CometBroadcastHashJoinExec.
  • A sort-merge join on decimal(38,18) keys between a converted input and one that cannot convert, which keeps all its rows because both shuffles stay columnar.
  • Decimal, Option, nested struct and array<string> fields, and an array<int> column that declines with its fallback reason.
  • A typed operation at the top of the plan, where nothing is converted.
  • The user function running exactly once per row.
  • A limit, a mapPartitions function and code reading Dataset.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.
  • An exception thrown past the first batch failing the query with the original message.
  • The feature being off by default.

Every converted plan is also checked for row operators reading a columnar child without a transition, and for spurious ColumnarToRow fallback reasons. Dropping the transition insertion fails 5 tests, and dropping the ColumnarToRowTransition change 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 Dataset suites from the 4.1.3 tests jar, with Comet enabled through system properties: DatasetSuite, DatasetPrimitiveSuite, DatasetAggregatorSuite, DatasetOptimizationSuite, DatasetCacheSuite and DatasetSerializerRegistratorSuite, 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, and groupBy.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.

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.
@github-actions github-actions Bot added the enhancement New feature or request label Oct 3, 2026
@andygrove andygrove added run-benchmark-check Run the benchmark compile and lint check on this pull request instead of waiting for the merge queue run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue labels Oct 3, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
  • Design approach: Add an opt-in CometSparkToColumnarExec above supported SerializeFromObjectExec outputs. 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.1 at 366b157a7583f646ce2d6ca0f87595ebe807e11f, 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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[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().

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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.
@github-actions github-actions Bot added the area:shuffle Shuffle (JVM and native) label Oct 3, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
  • Design approach: Optionally insert CometSparkToColumnarExec above supported SerializeFromObjectExec outputs 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.1 at 4f5cf2db1e80bf94db8a4b79a870738d621f1262, 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.
@andygrove
andygrove requested a review from comphead October 3, 2026 17:24

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
  • Design approach: Optionally insert CometSparkToColumnarExec above supported SerializeFromObjectExec outputs 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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 sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
  • Design approach: Optionally convert SerializeFromObjectExec output 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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 sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
  • Design approach: Optionally convert SerializeFromObjectExec output 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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 sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native partial aggregation and joins.
  • Design approach: Optionally convert SerializeFromObjectExec output 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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 mapPartitions case.

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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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 sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
  • Design approach: Optionally convert SerializeFromObjectExec output to Arrow while retaining Spark execution of user functions.
  • Correctness / compatibility analysis: The earlier decimal-join, strict-compilation, physical-limit and mapPartitions findings 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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-map case.

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")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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.

@andygrove

Copy link
Copy Markdown
Member Author

@copilot resolve the merge conflicts on this branch.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
  • Design approach: Optionally convert SerializeFromObjectExec output 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) returns 2 in Spark and conversion-disabled Comet, but conversion-enabled Comet evaluates the throwing row 30 and fails. Earlier decimal-join, physical-limit, mapPartitions and 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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-map reproduction.

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 sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
  • Design approach: Optionally convert supported SerializeFromObjectExec output 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) returns 2 in Spark and conversion-disabled Comet, but conversion-enabled Comet evaluates the throwing row 30 and fails. Earlier decimal-join, physical-limit, mapPartitions and 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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-map reproduction.

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.
@andygrove
andygrove requested a review from sunchao October 5, 2026 14:02

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
  • Design approach: Optionally convert supported SerializeFromObjectExec output 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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 sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Operators above typed Dataset transformations stayed on Spark until the next shuffle, preventing native aggregation and joins.
  • Design approach: Optionally convert supported SerializeFromObjectExec output 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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.

@andygrove
andygrove added this pull request to the merge queue Oct 5, 2026

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

Comment on lines +1331 to +1334
// 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")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

#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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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?

Comment on lines +83 to +85
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.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

@andygrove

Copy link
Copy Markdown
Member Author

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 spark.comet.convert.rdd.enabled or spark.comet.exec.localTableScan.enabled, a struct column and a repartition, and the same code is in every release since 1.0.0. Native shuffle wraps the batches of any child that isn't a native operator in ColumnarBatchArrowReader, which closes the vectors the conversion reuses. It also fails at the default batch size once a partition has more than one batch. #6607 reads a CometNativeArrowSource child as an Arrow stream instead. With just that change applied on top of this branch, your repro, the orderBy and join variants, and the two leaf conversions all pass. Since this is already queued and the config is off by default, I'd rather let it land and add your struct-under-shuffle test to #6607, which will close #6685, when I merge main into it. A struct decline here would only come out again in #6607.

On the profiles, I ran CometTypedDatasetSuite against main plus this branch on 3.4, 4.0 and 4.2, and all 15 tests pass. With the earlier 3.5 and 4.1 runs, it's green on every profile. You're right about the top-k arm: SerializeFromObjectExec doesn't report an output ordering in any supported version, so I'll remove it. I'll do that, the #6005 TODO, the shared benchmark helper and the datasources.md bullet in #6607, since it touches the same code and lands next.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:shuffle Shuffle (JVM and native) enhancement New feature or request run-benchmark-check Run the benchmark compile and lint check on this pull request instead of waiting for the merge queue run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Fuse the SerializeFromObject / MapElements / DeserializeToObject sandwich into a Comet projection instead of falling back

3 participants