Skip to content

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

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

Visorgood wants to merge 11 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 – the rejection is what keeps the whole stage off Comet. CometCollationSuite treats that as the primary line of defense for collated keys and pins the reasons these checks emit (#1947). Since #6206 it is no longer the only guard: supportedSortType declines a non-default collation anywhere in a sort key, and the aggregate serdes decline one too. The checks also cover hash and range partitioning only, not SinglePartition or round robin.

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

CometUnserializableShuffleKeyBenchmark measures the exchange shapes discussed in the thread; results are in the comment below.

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. Its fields are nullable, so canHoldNestedNull holds, and with the default null ordering CometSortOrder falls through to the strict floating-point case and reports Incompatible (Native sort orders null elements of array and struct keys by the key's null order, unlike Spark #6476).
  • 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.

Verified on Spark 4.1 with scalastyle and spotless enabled:

CometShuffleSuite, DisableAQECometShuffleSuite, CometShuffleManagerSuite,
CometCollationSuite, CometExpressionSuite, CometTPCDSV1_4/V2_7_PlanStabilitySuite
Suites: completed 7, aborted 0
Tests: succeeded 439, 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.

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.

Good catch. Added item 5 to "When Native Shuffle is Used": native shuffle serializes the partitioning into its protobuf plan, so every HashPartitioning expression and every RangePartitioning sort order has to convert through QueryPlanSerde.exprToProto.

For the mapsort example I named CometMapSort explicitly, since item 4 mentions mapsort as what admits map keys and the two could otherwise look contradictory: the type is admissible with nested hash keys enabled, but CometMapSort supports scalar map keys only, so a map with array or struct keys fails the expression check. That matches the comment already in supportedHashPartitioningDataType. Cross-linked to jvm_shuffle.md so the two lists agree.

The workflows will need approving again.

@andygrove
andygrove enabled auto-merge September 29, 2026 00:00
@andygrove
andygrove added this pull request to the merge queue Sep 29, 2026
@github-merge-queue
github-merge-queue Bot removed this pull request from the merge queue due to failed status checks Sep 29, 2026
@Visorgood

Copy link
Copy Markdown
Contributor Author

@andygrove there was a CI flake. The only real failure was Iceberg Spark SQL Tests (Iceberg 1.11) / Build Native Library, and it died at the "Setup Rust & Java toolchain" step – before anything was compiled. The Iceberg shards were
then skipped, so iceberg-spark-shard-coverage found no inventories and failed too, which took Required Checks down with it. No Iceberg test actually ran.

PR #5876 queued on the same base commit (8369bf1) right after and passed the identical Iceberg 1.11 pipeline, native build included.

The head is clean and both approvals still stand – could you re-add it to the queue?

@Visorgood
Visorgood requested a review from andygrove September 29, 2026 17:32

@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, and 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 code 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.

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

Reviewed full SHA 57e637ad2d3d9ab3d6866946600213f65ddc6370 against base 36ab57c689d4551b060637ec311f0e8da3b96a62, covering all six changed files and eight 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: 23 checks passed, 14 skipped, none failed. Comet CI includes passing Required Checks, TPC-DS/TPC-H verification and Celeborn compatibility checks. Logs confirm 530 shuffle tests and 1,544 expression tests passed, including the new regressions. These jobs tested GitHub merge commit 55eb81d05a79b27c291fc3254ecb72759598d3a7, which includes newer main changes. Spark SQL suites were skipped at this head. All seven Spark SQL 4.1 shards passed at predecessor d620618502c95e7b389930d3c9cc36db6c6abcb5, whose tested merge tree matches that predecessor exactly. The only subsequent change is documentation.

The separate merge-queue run failed because the Rust toolchain download reported Network is unreachable, before compilation. Iceberg coverage and required checks consequently failed.

Validation limits: upstream source comparisons and git diff --check passed. No local JVM/native build or runtime suites ran. This checkout lacks compiled artifacts and cached Spark dependencies. Runtime evidence comes from Spark 4.1 CI. Other supported versions received source comparison only.

@andygrove

Copy link
Copy Markdown
Member

I would like to do a deeper review of this PR before it is merged.

@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 and projections match upstream Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Writer routing keeps JVM dependencies separate from native serialization. Downstream native sorts retain their compatibility checks.
  • Key design decisions: Preserve native serialization requirements, payload-type checks, collation guards and Celeborn restrictions. Removing redundant planning checks simplifies the code without adding an abstraction or runtime algorithm. Performance was not benchmarked.
  • Implementation sketch: Tests cover exchange selection, per-row hash partition assignments and collation fallback with AQE enabled and disabled. Documentation describes 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 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.

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

Reviewed full SHA 57e637ad2d3d9ab3d6866946600213f65ddc6370 against base 36ab57c689d4551b060637ec311f0e8da3b96a62, covering all six changed files and eight base-relative commits. Confirmed the PR remains non-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: 23 checks passed, 14 skipped, none failed. Comet CI passed required checks, TPC-DS/TPC-H verification and Celeborn compatibility checks. Inspected logs show 530 shuffle tests and 1,544 expression tests passed, including the regressions. These jobs tested merge commit 55eb81d05a79b27c291fc3254ecb72759598d3a7, which includes newer main changes. Spark SQL suites were skipped at this head. All seven Spark SQL 4.1 shards passed at predecessor d620618502c95e7b389930d3c9cc36db6c6abcb5. Its tested merge tree matches that predecessor, and the only subsequent change is documentation.

The separate merge-queue run failed downloading the Rust toolchain with Network is unreachable, before compilation. Iceberg coverage and required checks consequently failed.

Validation limits: upstream source verification, cross-version comparisons and git diff --check passed. No local JVM/native suites ran. This checkout lacks compiled artifacts and cached Spark dependencies. Runtime evidence comes from Spark 4.1 CI; other supported versions received source comparison only.

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

I went back over this against branch-1.1, since that is what users will upgrade from, as well as against current main. The core change holds up. On the columnar path the partition ids and range bounds come from Spark's own partitionIdExpression and RangePartitioner, which match Spark from 3.4.3 through 4.2.0, so moving these exchanges off Spark's shuffle can't change where rows land, and your new assignment tests pin exactly that. Below is one request that goes with the merge from main you'll need for the conflict, and one question about cost when the operator reading the exchange stays on Spark.

1. spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffleExchangeExec.scala:722

This branch conflicts with main in CometExpressionSuite now. main rewrote the two nested floating-point sort tests around checkStrictNestedFloatingPointSort, which expects can hold a null element or field, so I'd take main's side of both hunks. That reason is recorded on the Sort's SortOrder, so the tests should still find it with the exchange on columnar shuffle.

The same merge brings in #6206, which closed #6158. supportedSortType now declines a non-default collation in any sort key at any depth, so a few of the new comments describe a gap that main no longer has:

  • the comments on the hash and range collation checks here, which say supportedSortType only type-checks single-column sorts and that the sort gap is tracked in #6158
  • the last sentence of the new collation paragraph in jvm_shuffle.md, which says this fallback is currently what stops CometSort from ordering such a key by raw bytes
  • the comment inside collation introduced above the scan still falls back to Spark's shuffle, which says declining the exchange is what keeps CometSort off the collated key

I'd still keep the checks. CometCollationSuite calls them the primary line of defense and pins their reasons, and with #6206 CometSort now declines collated keys on its own, as the hash aggregate and join serdes already did. Could the comments say that instead? Otherwise the next reader is back to wondering whether the checks are dead code, which is where I started.

2. General

I have a question about cost. For ORDER BY, joins and windows on one of these keys, the operator above the exchange serializes the same key again (CometSortExec, the join serdes and CometWindowExec all do), so it fails the same check and stays on Spark. For those shapes this PR swaps Spark's shuffle for the JVM columnar shuffle with a Spark operator reading it.

With a native child, the 1.1 plan is Comet's columnar-to-row transition, then Spark's Exchange, then the Sort. With this PR it is CometColumnarExchange, which pulls rows through Spark's ColumnarToRowExec and converts them back to Arrow, and then another columnar-to-row transition in front of the Sort. With a non-Comet child, which spark.comet.shuffle.convertFromSparkPlan.enabled admits by default, no Comet operator sits on either side. That is the round trip revertRedundantColumnarShuffle removes for aggregate pairs (#4004). These shuffles also move from Spark's execution memory to Comet's off-heap pool.

I don't expect a big number, since I only saw about 2% on TPC-DS when I measured the aggregate case for #4010, but nothing has been measured for this change. Could you time one of these against branch-1.1? In this snippet, turning the dispatcher off only stands in for a key without a serde, and a Hive UDF would do the same.

withSQLConf(CometConf.COMET_SCALA_UDF_CODEGEN_ENABLED.key -> "false") {
val bump = udf((x: Int) => x + 1)
spark.read.parquet(path)
.orderBy(bump(col("a")), col("b"))
.write.format("noop").mode("overwrite").save()
}

If it comes out slower, one option is to keep the probe when the child is not a Comet plan. That avoids the case with Spark on both sides and still lets a native child feed the columnar shuffle.

@Visorgood

Copy link
Copy Markdown
Contributor Author

Rebased onto main and addressed both points.

The CometExpressionSuite conflict took main's side of both hunks, as you suggested – the file is now identical to main and has dropped out of this PR. I also found a fourth stale comment you hadn't listed: the one above range partitioning on a nested floating-point key claimed CometSortOrder declines the key "for nested floating point under strict mode". After #6476 that is no longer the rule – canHoldNestedNull has to hold first. The test passes only because the parquet struct's fields are nullable, and the comment now says so.

On cost: four shapes, two runs each, 8M rows on an M3 Pro. The baseline arm switches Comet shuffle off, which reproduces the pre-#5971 plan inside one build; comparing against branch-1.1 directly would also carry every unrelated commit between us.

                                                          run 1        run 2

range, Comet scan +5.7% +2.9%
range, row-based Spark scan -1.6% -1.5%
hash, Comet scan -11.7% -11.1%
hash, row-based Spark scan -4.2% -6.9%

One shape is consistently slower, three are consistently faster. The direction reproduces, the magnitude does not – run-to-run drift is wider than the within-run stdev, so I'd call it 1.5-6% rather than pin a number. That is close to the ~2% you saw on the aggregate case.

Two things worth flagging:

The slow shape is the one with a native child, not the one with Spark on both sides. Keeping the probe when the child is not a Comet plan would preserve the regression and give up the wins.

And Spark-on-both-sides is narrower than it looks: columnarShuffleFailureReasons takes a child only if it is row-based or a Comet operator, so a Spark columnar scan is declined outright. That case only arises with a row-based child, where this change is faster.

Added the harness as CometUnserializableShuffleKeyBenchmark so this can be re-run. Happy to put the probe back for a narrower case if you'd rather, but on these numbers I'd leave it.

@Visorgood

Copy link
Copy Markdown
Contributor Author

The run-spark-4.1-tests label came off at some point; worth re-adding given the merge from main.

@Visorgood
Visorgood requested a review from andygrove October 11, 2026 07:48

@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 keys that native serde could not serialize, although Spark evaluates those keys on the JVM.
  • Design approach: Remove the hash and range exprToProto probes from columnarShuffleFailureReasons.
  • Correctness: No introduced partitioning-correctness defect identified. Hash projections, range sampling and range ordering match Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The new hash benchmark has a reproducible validation defect described in the finding.
  • Compatibility analysis: Native serialization requirements, payload-type checks, collation guards and Celeborn restrictions remain. Spark evaluates mapsort and unsupported native expressions itself. Downstream native operators retain their compatibility checks.
  • Key design decisions: Keep JVM partitioning eligibility separate from native expression support while preserving existing fallback boundaries.
  • Implementation sketch: Two probes are removed. Six new tests cover exchange selection, hash partition assignments and collation fallback in both AQE configurations. Three documentation pages and a four-case benchmark describe and exercise the expanded routing.
  • Performance: Removing unused planning probes avoids serialization work. Newly admitted exchanges can add row/Arrow conversions. The author's two runs report the native-scan range case becoming 5.7% and 2.9% slower. This previously raised performance concern remains unresolved. The reported hash improvements cannot establish this PR's benefit because that benchmark's exchange key already passes the old gate.
  • Design: Reusing Spark's partitioning implementation makes the eligibility change straightforward. The existing concern about unnecessary conversions above a native scan still needs resolution.
  • Abstraction & complexity: No production abstraction or partitioning algorithm is added. The benchmark's case/config helpers are reasonable, but checking exchange names alone does not establish that it exercises the removed gate.
  • Behavioral changes worth calling out: Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, unsupported native expressions and nullable nested-floating-point range keys can now use JVM columnar shuffle. This intentionally changes routing. The measured range slowdown is an accompanying regression, not a correctness improvement.
  • Suggested improvements: Repair the hash benchmark to retain the unsupported expression inside the exchange key, verify that the old probe rejects it, and rerun the comparison. Address the existing native-scan range slowdown before merging.

Reviewed all six files in the full diff from 00a4b422f0ed6560ef76dce004b94c3697613108 to aa00ed0353844322ab6dc08f3036103fe0dc2b1c, including all prerequisite commits. Confirmed the PR remains non-draft. Read the existing reviews, conversation 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 title checks report action_required. There is no exact-head build/test verdict, and run-spark-4.1-tests is absent.

Validation: git diff --check passed. A disposable Spark 4.1.3 probe compiled and executed the benchmark queries and a direct-hash control with AQE enabled and disabled, confirming the finding. Upstream source comparisons covered all supported profiles. Exact-head Comet JVM/native suites and throughput benchmarks were not run because this checkout has no compiled Comet artifacts. The author's test counts and timing measurements were not independently reproduced. No project code or GitHub state was changed.

/** Range partitioning, as in the reported shape, and hash partitioning for comparison. */
private val rangeQuery = "SELECT * FROM parquetV1Table ORDER BY bump(a), b"
private val hashQuery =
"SELECT max(b) FROM parquetV1Table GROUP BY bump(a) ORDER BY 1 LIMIT 10"

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] Use a hash query whose exchange key actually fails the removed serde probe. Spark computes bump(a) below the partial aggregate here, so the exchange hashes the integer attribute _groupingexpression, which already passes exprToProto on the base revision. Consequently, disabling Comet shuffle does not reproduce this query's pre-PR plan, and the reported hash speedups do not measure this change. checkPlans misses this because it only checks exchange names. This invalidates two of the comparisons used to weigh the existing range slowdown. Use a direct repartition(..., bump(col("a"))) or DISTRIBUTE BY bump(a), assert that the exchange key fails native serde with the dispatcher disabled, and rerun the comparison.

Evidence: A compiled Spark 4.1.3 probe over 100 Parquet rows executed the exact benchmark SQL with AQE=false and AQE=true. Both produced hashpartitioning(_groupingexpression#..., 4) with UDF_IN_EXCHANGE=false. The control SELECT * FROM parquetV1Table DISTRIBUTE BY bump(a) produced a hash key containing bump with UDF_IN_EXCHANGE=true. Spark's AggUtils.planAggregateWithoutDistinct uses grouping attributes for the final aggregate's distribution. Comet's unchanged CometAttributeReference serde accepts this integer attribute, so the old columnar probe does not reject it. Reproduction source and output: /tmp/review6110-planprobe-6f4_q1ri/Review6110PlanProbe.scala and /tmp/review6110-planprobe-6f4_q1ri/probe.log.

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