Skip to content

fix: stop gating columnar shuffle on native-serde checks it never uses - #6110

Open
Visorgood wants to merge 7 commits into
apache:mainfrom
Visorgood:visorgood/5971-columnar-range-partitioning
Open

Visorgood wants to merge 7 commits into
apache:mainfrom
Visorgood:visorgood/5971-columnar-range-partitioning

Conversation

@Visorgood

@Visorgood Visorgood commented Sep 22, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #5971.

Rationale for this change

columnarShuffleFailureReasons asked whether Comet could serialize the partitioning expressions to protobuf, in both the RangePartitioning and HashPartitioning branches. Nothing on the columnar path consumes that:

  1. prepareJVMShuffleDependency partitions on the JVM – UnsafeProjection over h.partitionIdExpression / the sort keys, LazilyGeneratedOrdering, Spark's RangePartitioner.
  2. CometShuffleDependency.outputPartitioning is a Catalyst Partitioning, not a proto message.
  3. PartitioningOuterClass.RangePartition is built only in CometNativeShuffleWriter.
  4. That writer is reached only through CometNativeShuffleHandle; CometShuffleManager hands the columnar path a different handle.
  5. CometCelebornShuffleManager also reads the partitioning, but only via nativeDependency, which requires shuffleType == CometNativeShuffle; rejectCometHandle throws for both columnar handles.

So both probes rejected exchanges the JVM would have partitioned correctly, and those queries fell back to Spark's shuffle for no compatibility reason. This is an unnecessary-fallback bug, not a correctness bug.

The structurally identical probes in the native branch are untouched – there the serialized expressions really do go native.

What changes are included in this PR?

Dropped the exprToProto probe from both branches of columnarShuffleFailureReasons.

The collation checks stay, and my earlier rationale for keeping them was wrong: they are reachable. CometScanRule only rejects a stored collated column, so _2 COLLATE UTF8_LCASE over a plain-string Parquet column keeps CometNativeScan native and the collation arrives in a Project above it. VALUES reaches the gate too.

They also matter beyond the shuffle. Spark's collation-aware Murmur3Hash / LazilyGeneratedOrdering run on the JVM here, so no native step sees a collated partition key – but rejecting the exchange moves the whole stage off Comet, which is what keeps CometSort away from the collated key. supportedSortType only type-checks single-column sorts, so a multi-column collated sort slips through. With the checks removed, listagg DISTINCT ... COLLATE utf8_lcase returns aabb instead of ab and #1947 regresses. The checks cover hash and range partitioning only, so they are not a complete guard – @andygrove filed #6158 for the sort gap itself.

inputs and the QueryPlanSerde import are both still used in the method.

How are these changes tested?

Five new tests in CometColumnarShuffleSuite, each confirmed to fail before the change:

  • range partitioning on a nested floating-point key – a struct<double, int> sort key, which CometSortOrder reports Incompatible for under strictFloatingPoint because strictFloatingPointReason recurses through containsType.
  • range partitioning on an unserializable expression and hash partitioning on an unserializable expression – a Scala UDF with spark.comet.exec.scalaUDF.codegen.enabled=false, so CometScalaUDF.convert returns None. Turning the dispatcher off is just a stable way to get a partition key with no serde; unlike the two cases above, nothing here depends on strictFloatingPoint.
  • two partition assignment matches Spark tests comparing spark_partition_id() per row. checkShuffleAnswer only compares the answer, which is order-insensitive and would pass even if Comet routed rows to different partitions, so assignment needs its own check (the same reasoning as CometNativeShuffleSuite). Both also assert one CometShuffleExchangeExec in the Comet run, so a future fallback cannot leave both sides on plain Spark and pass while testing nothing.

One more test, collation introduced above the scan still falls back to Spark's shuffle, passes before and after: it pins the collation guard this PR keeps, which the suite did not cover.

One existing expectation changed: columnar shuffle on array/struct map key/value expected 0 Comet exchanges on Spark 4.0+, because Spark wraps map shuffle keys in mapsort(...) and Comet cannot serialize that for array or struct map keys. The columnar path computes partition ids on the JVM from h.partitionIdExpression, mapsort included, so that verdict never applied to it. The expectation is now 1, and the new map-key assignment test pins that the distribution still matches Spark.

CometExpressionSuite's two nested floating-point sort tests asserted the removed reason string; they now assert the CometSortOrder reason, which is what those tests are actually about – the Sort still falls back, only the exchange no longer does.

Verified on Spark 4.1 with scalastyle and spotless enabled:

CometShuffleSuite, DisableAQECometShuffleSuite, CometShuffleManagerSuite,
CometCollationSuite, CometTPCDSV1_4/V2_7_PlanStabilitySuite
Suites: completed 6, aborted 0
Tests: succeeded 256, failed 0

Both AQE configurations, since checkCometExchange strips the AQE plan. The plan-stability
goldens are unchanged, so nothing needed regenerating.

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

Now that exchanges that used to fail the probe can go columnar in auto mode, could some TPC-DS/TPC-H plans pick up a Comet exchange where they previously fell back? CI hasn't run the plan-stability suites yet. Once it does, if any goldens move, please regenerate them in this PR with dev/regenerate-golden-files.sh.

Also, #5802 changes the same columnar shuffle on array/struct map key/value test in the other direction. It keeps 0 exchanges on 4.0+ and adds a flag-gated columnar test. Once this lands, that test passes without the flag. Could you and @sam-1112 coordinate which lands first, so the other can drop or adjust its columnar test?

}
}
for (dt <- expressions.map(_.dataType).distinct) {
if (isStringCollationType(dt)) {

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.

Thanks for tracing this through so carefully. The rationale really helped. One question about the collation checks you kept. The argument for dropping the exprToProto probes is that partitionIdExpression and the range ordering run on the JVM through UnsafeProjection and LazilyGeneratedOrdering. Doesn't that apply to collated string keys too? Spark's own collation-aware hash and ordering would run there. If there's a native step on the columnar path that a collated partition key reaches, could you point to it in a comment? If not, I think these checks should go too, so the function doesn't keep a guard the PR's own reasoning says is unnecessary.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

You're right that nothing native sees the key: the partition id comes from Spark's Murmur3Hash / LazilyGeneratedOrdering on the JVM, and both are collation-aware.

But I believe the checks still can't go. Rejecting the exchange is what moves the whole stage off Comet, and that also keeps CometSort off the collated key. supportedSortType only type-checks single-column sorts (QueryPlanSerde.scala:1287), so a multi-column collated sort gets past it.

I removed both checks and ran CometCollationSuite: 4 failures. Three only change the reason string the test pins, since the query still falls back via the sort check. The fourth is a wrong answer:

The plan keeps CometColumnarExchange and CometSort over a two-column collated sort key, so Comet dedups a/A on raw bytes and #1947 is back.

Kept them, and rewrote the rationale in the description – my "unreachable" claim there was wrong. Also added a test for the fallback, which the suite didn't have. I'll file a separate issue for the single-column limit in supportedSortType; fixing that is what would make this check redundant.

* pass even if Comet routed rows to different partitions than Spark. Compare
* spark_partition_id() per row instead.
*/
private def checkPartitionAssignmentMatchesSpark(df: => DataFrame, clue: String): Unit = {

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.

This helper only compares the Comet run with the Spark run. If a future change makes these queries fall back, both sides are plain Spark and the test still passes. Could the helper also assert one CometShuffleExchangeExec in the Comet run, for example with checkCometExchange(df, 1, false)? The map-key assignment test especially has no other test pinning that exact query to Comet.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Agreed, that would have passed while testing nothing. Added checkCometExchange(df, 1, false) at the top of the helper, so both assignment tests now pin the Comet run to one exchange before comparing partition ids.

@sam-1112

Copy link
Copy Markdown
Contributor

@andygrove @Visorgood Thanks for flagging #5802.

For collation, I agree that JVM columnar shuffle has no native partition-id calculation. However, scan fallback does not make this guard unreachable: a non-default collation can be introduced above a normal scan or come from VALUES. The existing collation suite relies on the shuffle rule as a safety boundary against raw-byte Comet sort or aggregate behavior. I would keep the guard, but revise the rationale rather than describe it as unreachable.

For #5802, I suggest #6110 lands first. I will then rebase #5802 and update its columnar-shuffle tests and docs: with the dispatcher off, #6110 can use JVM columnar shuffle; with it on, #5802 can enable native shuffle.

@Visorgood

Copy link
Copy Markdown
Contributor Author

Hey @sam-1112 ! You're right. CometScanRule only rejects a stored collated column, so _2 COLLATE UTF8_LCASE over a plain string column leaves CometNativeScan in place and the collation lands in a project above it – the exchange does reach the gate. VALUES gets there the same way, and CometCollationSuite already covers that.

Rewrote the rationale in the description and added a test for the fallback. It also turns out the check matters beyond the shuffle: with it removed, listagg DISTINCT under utf8_lcase returns aabb instead of ab, because CometSort stays on a two-column collated key. Details in my reply to Andy.

#6110 first works for me, thanks for offering to rebase #5802.

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

The collation reasoning in your reply convinced me the two checks in columnarShuffleFailureReasons have to stay. The trouble is that the reason only lives in this thread, so the next person to read those checks will likely conclude they're dead code, as I did. Could you add a short comment at the checks saying they keep the stage off Comet so CometSort never sees a multi-column collated key that supportedSortType lets through?

Two of the new test comments also say something different from what the tests do. The collation test says Comet hashes raw bytes and would misroute rows, but on this path the partition id comes from Spark's collation-aware Murmur3Hash. And the comment above the UDF tests says they reproduce at default config, but they set spark.comet.exec.scalaUDF.codegen.enabled=false, which defaults to true. Could both say what's actually going on?

@Visorgood

Copy link
Copy Markdown
Contributor Author

All three fixed, thanks.

The checks now carry the reasoning inline: that this isn't a shuffle-correctness check, and that the fallback is what keeps CometSort off a collated key supportedSortType lets through.

You're right about the UDF comment – scalaUDF.codegen.enabled defaults to true, so those tests are not at default config. Reworded to say what they actually do: turning the dispatcher off is a stable way to get a partition key with no serde, and the point is that nothing there depends on strict mode. Fixed the same claim in the description.

Plan stability: both suites pass unchanged, so there is nothing to regenerate.

CometTPCDSV1_4_PlanStabilitySuite, CometTPCDSV2_7_PlanStabilitySuite
Suites: completed 2, aborted 0
Tests: succeeded 129, failed 0

Full run of everything this touches, including CometCollationSuite:

Suites: completed 6, aborted 0
Tests: succeeded 252, failed 0

Ordering with #5802 is settled with @sam-1112 – this one first, then he rebases.

@Visorgood
Visorgood requested a review from andygrove September 23, 2026 19:30

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

Thanks for the updates. The inline comments at the collation checks read well, and I'm happy with the code. I traced both directions I was worried about. The removed probe never rejected nested collated keys, so nothing new gets through there (and #6158 covers any depth). Downstream operators that now stay on Comet still do their own checks, like CometSortOrder for nested floating point. A few things before this goes to the queue:

docs/source/contributor-guide/jvm_shuffle.md, "When JVM Shuffle is Used"

Since native shuffle still rejects partition keys it can't serialize, those now land on JVM columnar shuffle instead of Spark's. Could you add that as a case in the list, with mapsort over array or struct map keys on Spark 4.0+ as an example? It would also help to say that collated hash and range keys go past both Comet paths to Spark's shuffle, since the new collation test pins exactly that.

docs/source/user-guide/latest/understanding-comet-plans.md, the CometColumnarExchange paragraph

This paragraph gives collated strings as an example of a key that falls back to CometColumnarExchange. Your new test collation introduced above the scan still falls back to Spark's shuffle shows they actually go to Spark's shuffle. Could you fix the example while you're here? A Scala UDF key or a mapsort key would fit better now.

Spark SQL tests

Since this changes which exchanges become Comet in auto mode, I'd like to run the Spark SQL tests before queueing. Spark's SQL suites assert exchange shapes, and I'd rather see those results here than in the merge queue. I'll add the run-spark-4.1-tests label.

@andygrove andygrove added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 23, 2026
@Visorgood

Visorgood commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor Author

Both docs updated.

jvm_shuffle.md: added a fourth case to "When JVM Shuffle is Used" – native shuffle serializes the partitioning expressions, so a key it has no serde for keeps the exchange off the native path while JVM shuffle takes it, with the Spark 4.0+ mapsort wrapper over an array or struct map key as the example. Followed by a note that a collated hash or range key is declined by both Comet paths and stays a plain Spark Exchange, which is what the new collation test pins.

understanding-comet-plans.md: collated strings were given as the example for CometColumnarExchange, which is wrong – they go to Spark's shuffle. Replaced with the mapsort key, and said explicitly that a collated key is not one of these cases.

Thanks for adding run-spark-4.1-tests.

The workflows are still awaiting approval – the head commit has only one check run, so the Spark SQL suites haven't started yet.

@Visorgood
Visorgood requested a review from andygrove September 25, 2026 11:12

@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: JVM columnar shuffle rejected partition expressions that Comet could not serialize, although Spark evaluates those expressions on the JVM.
  • Design approach: Remove the hash and range exprToProto probes from columnarShuffleFailureReasons.
  • Correctness / compatibility analysis: Traced partition evaluation, dependency construction, writer selection and Celeborn routing. The relevant hash projections and complete range sampling/comparator setup match Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Downstream native sorts retain their own compatibility checks.
  • Key design decisions: Native serialization checks, payload-type checks and collation guards remain. Removing the unused probes simplifies planning without adding an abstraction or runtime algorithm. Performance was not benchmarked.
  • Implementation sketch: Tests assert columnar exchange selection, compare per-row partition IDs with Spark, and retain collation fallback coverage. Documentation explains the expanded JVM fallback path.
  • Behavioral changes worth calling out: Unsupported native expressions, complex-key mapsort expressions and strict nested-floating-point range keys can now use JVM columnar shuffle. Collated string hash/range keys still fall back to Spark.
  • Suggested improvements: No additional P1/P2 changes requested. Existing review concerns are addressed at this head.

No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 concerns remain unresolved.

Reviewed full SHA a7b75e122f1ce53f9e6f57a47397b019b3d0cffe, covering all five files and six commits in the full PR merge-base diff against supplied base 9a4d5f28368d4fbb54ac23c14c5d52d5637ce26d. Confirmed the PR is not a draft. Read the supplied discussion, reviews, inline comments and threads. Routed skills: review-comet-pr, review-comet-shuffle-pr, and review-comet-expression-pr.

Exact-head CI: only the label check passed. Comet CI and CodeQL report action_required. The run-spark-4.1-tests label is present, but there is no exact-head test verdict.

Validation limits: source comparisons and git diff --check passed. JVM suites could not run locally: compiled artifacts and Spark dependencies are absent, and Maven bootstrap with a writable temporary cache failed with UnknownHostException: repo.maven.apache.org. The author's reported test results were not independently reproduced. Runtime validation remains outstanding before queueing.

@Visorgood

Copy link
Copy Markdown
Contributor Author

Thanks for the thorough pass – especially for checking the hash projections and the range sampling/comparator setup against 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Cross-version behaviour was the one part I could only reason about rather than test.

Agreed on the CI point. The head has only the label check; Comet CI and CodeQL are both action_required. @andygrove added run-spark-4.1-tests, but the workflows still need someone to approve them before anything runs. Could one of you kick them off?

Everything I ran locally is listed in the description, and I'm happy to re-run or extend it once CI produces a verdict.

@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: JVM columnar shuffle rejected partition expressions that native serde could not serialize, although Spark evaluates those expressions on the JVM.
  • Design approach: Remove the hash and range exprToProto probes from columnarShuffleFailureReasons.
  • Correctness / compatibility analysis: Hash projections and range sampling/comparator setup match Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Native shuffle still requires serialization. Downstream native sorts retain their own compatibility checks.
  • Key design decisions: Preserve payload-type checks, collation guards and Celeborn routing restrictions. Removing unused planning probes simplifies the implementation without introducing a new runtime algorithm or abstraction. Throughput was not benchmarked.
  • Implementation sketch: Tests cover exchange selection, per-row hash partition assignments and collation fallback under both AQE configurations. Documentation explains the expanded JVM shuffle path.
  • Behavioral changes worth calling out: Expressions without native serde, complex-key mapsort expressions and nested-floating-point range keys can now use JVM columnar shuffle. Direct collated-string hash/range keys still fall back to Spark.
  • Suggested improvements: No additional P1/P2 changes requested. Existing review concerns are addressed in the current code and documentation.

No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 concerns remain unresolved.

Reviewed full SHA a7b75e122f1ce53f9e6f57a47397b019b3d0cffe, covering all five files and six commits in the full PR merge-base diff against supplied base 9a4d5f28368d4fbb54ac23c14c5d52d5637ce26d. Confirmed the PR is not a draft. Read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-expression-pr.

Exact-head CI: only the label check passed. Comet CI and CodeQL report action_required. The run-spark-4.1-tests label is present, but no exact-head test verdict is available.

Validation limits: upstream source verification, cross-version source comparisons and git diff --check passed. No JVM suites ran locally: compiled artifacts and Spark dependencies are absent, and a fresh Maven bootstrap failed with UnknownHostException: repo.maven.apache.org. The author's reported test results were not independently reproduced. Runtime validation remains outstanding before queueing.

…mnar-range-partitioning

# Conflicts:
#	docs/source/user-guide/latest/understanding-comet-plans.md
@Visorgood

Copy link
Copy Markdown
Contributor Author

Merged main to resolve the conflict. It was in the CometColumnarExchange paragraph of understanding-comet-plans.md: the 1.1.0 docs update (#6168) had already dropped collated strings as the example and added that collated keys use Spark's shuffle. I kept that wording and only added the case this PR introduces – native shuffle declining a key expression it cannot serialize, with the mapsort wrapper as the example.

Re-ran everything after the merge and a native rebuild:

CometShuffleSuite, DisableAQECometShuffleSuite, CometShuffleManagerSuite,
CometCollationSuite, CometTPCDSV1_4/V2_7_PlanStabilitySuite
Suites: completed 6, aborted 0
Tests: succeeded 256, failed 0

Plan-stability goldens are still unchanged.

The workflows are still awaiting approval, so there's no exact-head CI verdict yet.

@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: JVM columnar shuffle rejected partition expressions that native serde could not serialize, although Spark evaluates those expressions on the JVM.
  • Design approach: Remove the hash and range exprToProto probes from columnarShuffleFailureReasons.
  • Correctness / compatibility analysis: Hash projections and range sampling/comparator setup match upstream Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Writer selection keeps columnar dependencies separate from native serialization. Downstream native sorts retain their expression compatibility checks.
  • Key design decisions: Preserve native serialization requirements, payload-type checks, collation guards and Celeborn restrictions. Removing redundant planning checks simplifies the implementation without introducing an algorithm or abstraction. Throughput was not benchmarked.
  • Implementation sketch: Tests cover exchange selection, per-row hash partition assignments and collation fallback under both AQE configurations. Documentation explains the expanded JVM shuffle path.
  • Behavioral changes worth calling out: Expressions without native serde, complex-key mapsort expressions and strict nested-floating-point range keys can now use JVM columnar shuffle. Direct collated-string hash/range keys still fall back to Spark.
  • Suggested improvements: No additional P1/P2 changes requested. Existing review concerns are addressed in the current code and documentation.

No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 concerns remain unresolved.

Reviewed full SHA d620618502c95e7b389930d3c9cc36db6c6abcb5 against base 36ab57c689d4551b060637ec311f0e8da3b96a62, covering all five changed files and seven base-relative commits. Confirmed the PR is not a draft. Read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-expression-pr.

Exact-head CI: the label check passed. Comet CI, CodeQL and Check PR Title report action_required. The run-spark-4.1-tests label is present, but no exact-head test verdict is available.

Validation limits: upstream source verification, cross-version source comparisons and git diff --check passed. No JVM suites ran locally. Compiled artifacts and cached Spark dependencies are absent, and fresh Maven bootstrap failed with UnknownHostException: repo.maven.apache.org. The author's reported 256 passing tests were not independently reproduced. Runtime validation remains outstanding before queueing.

@Visorgood
Visorgood requested a review from sunchao September 26, 2026 11:20

@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: JVM columnar shuffle rejected partition expressions that native serde could not serialize, although Spark evaluates those expressions on the JVM.
  • Design approach: Remove the hash and range exprToProto probes from columnarShuffleFailureReasons.
  • Correctness / compatibility analysis: Hash projections, range projections, and range sampling/ordering setup match upstream Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Writer selection keeps columnar dependencies separate from native serialization. Downstream native sorts retain their expression compatibility checks.
  • Key design decisions: Preserve native serialization requirements, payload-type checks, collation guards and Celeborn restrictions. Removing redundant planning checks simplifies the implementation without introducing a runtime algorithm or abstraction. Performance was not benchmarked.
  • Implementation sketch: Tests cover exchange selection, per-row hash partition assignments and collation fallback with AQE enabled and disabled. Documentation explains the expanded JVM shuffle path.
  • Behavioral changes worth calling out: Expressions without native serde, complex-key mapsort expressions and strict nested-floating-point range keys can now use JVM columnar shuffle. Direct collated-string hash/range keys still fall back to Spark.
  • Suggested improvements: No additional P1/P2 changes requested. Existing review concerns are addressed in the current code and documentation.

No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 concerns remain unresolved.

Reviewed full SHA d620618502c95e7b389930d3c9cc36db6c6abcb5 against base 36ab57c689d4551b060637ec311f0e8da3b96a62, covering all five changed files and seven base-relative commits. Confirmed the PR is not a draft. Read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-expression-pr.

Exact-head CI: the label check passed. Comet CI, CodeQL and Check PR Title report action_required. The run-spark-4.1-tests label is present, but no exact-head test verdict is available.

Validation limits: upstream source verification, cross-version source comparisons and git diff --check passed. No JVM suites ran locally. Compiled artifacts and cached Spark dependencies are absent, and fresh Maven bootstrap failed with UnknownHostException: repo.maven.apache.org. The author's reported 256 passing tests were not independently reproduced. Runtime validation remains outstanding before queueing.

@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: JVM columnar shuffle rejected partition expressions that native serde could not serialize, although Spark evaluates those expressions on the JVM.
  • Design approach: Remove the hash and range exprToProto probes from columnarShuffleFailureReasons.
  • Correctness / compatibility analysis: Hash projections, range projections, and range sampling/ordering setup match upstream Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Writer selection keeps columnar dependencies separate from native serialization. Downstream native sorts retain their expression compatibility checks.
  • Key design decisions: Preserve native serialization requirements, payload-type checks, collation guards and Celeborn restrictions. Removing redundant planning checks simplifies the implementation without adding an abstraction or per-row work. Performance was not benchmarked.
  • Implementation sketch: Tests cover exchange selection, per-row hash partition assignments and collation fallback with AQE enabled and disabled. Documentation explains the expanded JVM shuffle path.
  • Behavioral changes worth calling out: Expressions without native serde, complex-key mapsort expressions and strict nested-floating-point range keys can now use JVM columnar shuffle. Direct collated-string hash/range keys still fall back to Spark.
  • Suggested improvements: No additional P1/P2 changes requested. Existing review concerns are addressed in the current code and documentation.

No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 concerns remain unresolved.

Reviewed full SHA d620618502c95e7b389930d3c9cc36db6c6abcb5 against base 36ab57c689d4551b060637ec311f0e8da3b96a62, covering all five changed files and seven base-relative commits. Confirmed the PR is not a draft. Read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-expression-pr.

Exact-head CI: the label check passed. Comet CI, CodeQL and Check PR Title report action_required. The run-spark-4.1-tests label is present, but no exact-head test verdict is available.

Validation limits: upstream source verification, cross-version source comparisons and git diff --check passed. No JVM suites ran locally. Compiled artifacts and cached Spark dependencies are absent, and a fresh Maven bootstrap failed with UnknownHostException: repo.maven.apache.org. The author's reported 256 passing tests were not independently reproduced. Runtime validation remains outstanding before queueing.

@Visorgood

Copy link
Copy Markdown
Contributor Author

CI is green on the current head: 32 checks passed, none failed, and Required Checks is green. The run-spark-4.1-tests label did its job – all seven Spark SQL 4.1 shards (catalyst, sql_core-1..3, sql_hive-1..3) passed, along with Verify TPC-DS/TPC-H Results and the Celeborn reflection compatibility checks.

@andygrove this should be ready now.

@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: JVM columnar shuffle rejected partition expressions that native serde could not serialize, although Spark evaluates those expressions on the JVM.
  • Design approach: Remove the hash and range exprToProto probes from columnarShuffleFailureReasons.
  • Correctness / compatibility analysis: Hash projections and range sampling/ordering match upstream Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Native shuffle retains serialization checks. Downstream native sorts retain their expression compatibility checks.
  • Key design decisions: Preserve payload-type checks, collation guards and Celeborn restrictions. Removing redundant planning probes simplifies the implementation without adding an abstraction or per-row work. Performance was not benchmarked.
  • Implementation sketch: Tests cover exchange selection, per-row hash partition assignments and collation fallback with AQE enabled and disabled. Documentation explains the expanded JVM shuffle path.
  • Behavioral changes worth calling out: Expressions without native serde, complex-key mapsort expressions and strict nested-floating-point range keys can now use JVM columnar shuffle. Direct collated-string hash/range keys still fall back to Spark.
  • Suggested improvements: No additional P1/P2 changes requested. Existing review concerns are addressed, and the requested Spark SQL validation now passes.

No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 concerns remain unresolved.

Reviewed full SHA d620618502c95e7b389930d3c9cc36db6c6abcb5 against base 36ab57c689d4551b060637ec311f0e8da3b96a62, covering all five changed files and seven base-relative commits. Confirmed the PR is not a draft. Read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-expression-pr.

Exact-head CI: 32 checks passed, 14 skipped, none failed. Comet CI includes passing Required Checks, all seven Spark SQL 4.1 shards, TPC-DS/TPC-H verification and Celeborn compatibility checks. Inspected logs confirm 530 passing shuffle tests, including the new regressions, and both changed expression tests passing. CI's merge commit has the same tree as the reviewed head.

Validation limits: upstream source comparisons and git diff --check passed. No local JVM/native build or runtime suites were run. Runtime evidence comes from Spark 4.1 CI. Other supported Spark versions were checked through source comparisons, without runtime validation.

@andygrove andygrove removed the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 27, 2026
Comment on lines +57 to +64
4. **Partition keys native shuffle cannot serialize**: native shuffle serializes the
partitioning expressions to protobuf, so a key expression Comet has no serde for, or whose
serde reports it incompatible, keeps the exchange off the native path. JVM shuffle has no
such requirement, because it evaluates the key on the JVM through `UnsafeProjection` and
`LazilyGeneratedOrdering`, so these exchanges land here rather than on Spark's shuffle. One
example is the `mapsort(...)` wrapper Spark 4.0 and later adds around a map used as a
shuffle key: Comet cannot serialize it for array or struct map keys, so such an exchange
becomes `CometColumnarExchange`.

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.

Thanks for the doc updates. One gap I noticed. This new case 4 says native shuffle declines a key expression Comet can't serialize. The "When Native Shuffle is Used" list in native_shuffle.md says native is chosen "when all of the following conditions are met", but it has no such condition. Could you add a fifth item there saying every hash partitioning expression and range sort order has to convert through exprToProto, with the same mapsort example? Then the two docs agree on when native is picked.

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) bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Columnar shuffle rejects range partitioning based on a native-serde check it never uses

4 participants