Skip to content

fix: coalesce shuffle partitions under Comet unions the way Spark does - #6459

Open
andygrove wants to merge 1 commit into
apache:mainfrom
andygrove:fix/cometunion-coalesce-partitions
Open

andygrove wants to merge 1 commit into
apache:mainfrom
andygrove:fix/cometunion-coalesce-partitions

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6454.

Rationale for this change

With AQE on, a union that Comet converts to CometUnionExec got no shuffle partition coalescing when one of its branches has a leaf that is not a shuffle, such as a scan or a table-cache stage, so the shuffled branches kept spark.sql.shuffle.partitions partitions. The issue has the analysis and a reproduction that ends with 201 partitions where Spark has 2.

Spark's CoalesceShufflePartitions coalesces each child of a UnionExec as its own group, and from Spark 4.0 each child of a CartesianProductExec, BroadcastHashJoinExec or BroadcastNestedLoopJoinExec too, but it matches those classes. Comet replaces them while AQE prepares a stage, before the optimizer rules run, so Spark's rule sees the Comet operators and falls through to the case that coalesces only when every leaf below is an exchange stage.

What changes are included in this PR?

A new query-stage optimizer rule, CometCoalesceShufflePartitions, registered after Comet's other two, so it runs after Spark's CoalesceShufflePartitions. It looks for Comet operators whose originalPlan is one of the operators above, and whose shuffle stages no AQE rule has put a read over. For each one it rebuilds the Spark operators over the Comet children, runs Spark's own CoalesceShufflePartitions over that, and swaps the Comet operators back in over the new children. Spark's code makes every coalescing decision, per Spark version.

The rule extends AQEShuffleReadRule, so AQE treats it like Spark's rule: it is skipped for the final stage when spark.sql.adaptive.applyFinalStageShuffleOptimizations is off, and its result is discarded if it breaks a distribution required above it.

It also covers a union whose branches are all shuffles with different partition counts. Spark's rule coalesces a Comet union's shuffles together, finds the counts differ, and coalesces nothing, where for a Spark union it would coalesce each branch on its own.

Two limits:

  • When every leaf below a Comet union is an exchange stage and the counts agree, Spark's rule already coalesces its shuffles, together rather than branch by branch, and this rule leaves them as they are.
  • Spark 3.4 has no hook for query-stage optimizer rules, so nothing changes there.

Once this lands, #5634 can drop the IgnoreComet it puts on Spark's SPARK-42101 union test.

How are these changes tested?

  • A new test in CometExecSuite runs the issue's query, a union of a repartitioned DataFrame with a Parquet scan. It checks for a coalesced read under the CometUnionExec, and for the same number of result partitions as with Comet disabled.
  • A port of Spark's SPARK-42101: Coalesce shuffle partition with union even if exists TableCacheQueryStage to CometInMemoryCacheSuite.

Both fail with the rule unregistered. They pass on Spark 3.5, 4.0, 4.1 and 4.2, and are canceled on 3.4.

Spark 4.1's AdaptiveQueryExecSuite, CoalesceShufflePartitionsSuite, DataFrameSetOperationsSuite, ExchangeSuite and DataFrameSuite, run from the test jar with Comet enabled, give the same results with and without the rule, except that CoalesceShufflePartitionsSuite's Union two datasets with different pre-shuffle partition number fails without it and passes with it.

CometExecSuite, CometJoinSuite, CometSetOpWithGroupBySuite, CometPlanEqualitySuite, CometInMemoryCacheSuite and CometNativeShuffleSuite pass on 4.1. The strict-warnings compile and scalafix pass on 3.5.

Spark's CoalesceShufflePartitions coalesces each child of a UnionExec as its
own group, and from Spark 4.0 each child of a CartesianProductExec,
BroadcastHashJoinExec or BroadcastNestedLoopJoinExec too. It matches those
classes, so a CometUnionExec fell through to the case that coalesces only when
every leaf below it is an exchange stage. A union with a scan or a table-cache
stage in one branch kept every partition of the shuffles in the others.

Add CometCoalesceShufflePartitions, a query-stage optimizer rule that runs
after Spark's. For a Comet operator whose Spark original is one of those, and
whose shuffle stages no AQE rule has touched, it rebuilds the Spark operators
over the Comet children, runs Spark's own rule over that, and swaps the Comet
operators back in. It extends AQEShuffleReadRule, so AQE gates and validates
it the way it does Spark's rule. Spark 3.4 has no hook for query-stage
optimizer rules, so the fix applies from Spark 3.5.

Closes apache#6454.
@github-actions github-actions Bot added the bug Something isn't working label Sep 30, 2026

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Thanks for the thorough write-up, and for reusing Spark's own rule instead of re-implementing the coalescing decisions. I left two questions inline about cases the rule might not cover or might over-cover, plus a few small nits. I haven't built or run anything, so the points about plan shapes are based on reading the code.

}
}
val coalesced = CoalesceShufflePartitions(SparkSession.active).apply(asSpark)
if (coalesced eq asSpark) plan else restore(coalesced)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

On Spark 4.1+ a union advertises its children's partitioning when they match (spark.sql.unionOutputPartitioning), so a final aggregate above it can skip its shuffle. If only one branch is coalescible, for example the other is repartition(200, $"k"), coalescing it alone changes that partitioning, and I expect AQE's ValidateRequirements to accept it because CometHashAggregateExec doesn't declare requiredChildDistribution. I haven't run this, so it may not be reachable. Would it make sense to skip operators whose outputPartitioning isn't UnknownPartitioning, with a 4.1 test that groups by the key over such a union?

val standIn = original.withNewChildren(p.children)
// `withNewChildren` hands back the original itself when the children are the same ones,
// and the tag must not land on the operator that the Comet one keeps.
if (standIn eq original) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

If a Comet union is created during an AQE re-plan on top of stages that are already materialized, I think originalPlan.children equals children, so withNewChildren returns original and this branch returns p without coalescing. I haven't run it, but AQE adopts a re-planned plan when it differs at equal cost, so I'd expect this to happen for queries where a join changes strategy below the union. Would it make sense to force a fresh copy of original here (for example with makeCopy) so those unions are handled too, and to cover that shape in a test?

}

// https://github.com/apache/datafusion-comet/issues/6454
test("AQE coalesces the shuffle partitions of a union whose other branch is a scan") {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

The description says the rule also covers a union whose shuffles Spark's rule gives up on together, but both new tests use a scan or a table cache stage as the other branch. It might be worth adding that shape, for example a hash-shuffled join unioned with a global aggregate as in Spark's Union two datasets with different pre-shuffle partition number. It differs from these tests because every leaf is an exchange stage, and I expect the SinglePartition shuffle of the aggregate to be what makes Spark's rule skip the whole Comet union.

private def replaced(plan: SparkPlan): Option[SparkPlan] = plan match {
case comet: CometExec =>
comet.originalPlan match {
case original @ (_: UnionExec | _: CartesianProductExec | _: BroadcastHashJoinExec |

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nit: I don't see a Comet counterpart of CartesianProductExec under spark/src/main, so that arm looks unreachable today. The BroadcastHashJoinExec and BroadcastNestedLoopJoinExec arms have no test either, and I'm not sure which plan shape triggers them, since a Comet join with only exchange stages below it is already coalesced by Spark's rule. Would it make sense to narrow this to what a test exercises?

* Coalesces the shuffle partitions below a Comet operator that Spark's CoalesceShufflePartitions
* coalesces child by child but does not recognize.
*
* Spark coalesces each child of a `UnionExec` as a group of its own, and from Spark 4.0 each

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nit: I think this holds from Spark 3.5 rather than 4.0. In v3.5.8, CoalesceShufflePartitions.collectCoalesceGroups already has a case for each of UnionExec, CartesianProductExec, BroadcastHashJoinExec and BroadcastNestedLoopJoinExec (lines 150 to 157), and v3.4.3 only has UnionExec. Might be worth adjusting here and in the PR description.

}
injectQueryStageOptimizerRuleShim(extensions, CometPlanAdaptiveDynamicPruningFilters)
injectQueryStageOptimizerRuleShim(extensions, CometReuseSubquery)
injectQueryStageOptimizerRuleShim(extensions, CometCoalesceShufflePartitions)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nit: the class scaladoc above lists the AQE query-stage optimizer rules for each stage and notes which ones are not registered on Spark 3.4. Would it make sense to add this rule to both places?

@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: Spark's AQE coalescer does not recognize Comet unions, leaving unnecessary shuffle partitions when another branch contains a scan or cache stage.
  • Design approach: Temporarily reconstruct Spark operators, reuse Spark's coalescing rule, then restore the Comet operators. This avoids duplicating Spark's partition-sizing logic.
  • Correctness / compatibility analysis: Compared Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Found one new P2 issue involving lost ancestor-join context. The existing partitioning concern also remains a substantiated correctness blocker: the planner probe changes union branches from [20,20] to [1,20], loses hash partitioning, and passes validation with Comet's unspecified aggregate distribution. A bounded reproduction using the production union RDD shim then produces duplicate (key,1) results instead of one (key,2) result per key. This existing concern is not duplicated in the findings.
  • Key design decisions: Extending AQEShuffleReadRule preserves final-stage optimization controls. However, invoking Spark's rule on isolated subtrees loses information Spark uses to preserve join parallelism.
  • Implementation sketch: Adds the tagged reconstruction/restoration rule, registers it after existing query-stage optimizer rules, and adds scan/union and cache/union tests.
  • Behavioral changes worth calling out: The intended scan/union case coalesces successfully. Existing AQE reads and the disabled-coalescing setting remain untouched. Spark 3.4 does not register the rule.
  • Suggested improvements: Preserve ancestor context when coalescing, resolve the existing union distribution blocker, and cover both reproduced cases with regression tests.

Reviewed full SHA bae0ffdba703df7e5810660fecce9d04183e3772 against base 9f68a4144fdefc70ec26f3d80b0e4eafe72a8df4. Inspected all 11 base-to-head file differences. Seven reflect main's newer #6366 fix absent from this branch and are unchanged from the common ancestor, so they are not reported as introduced regressions.

Routed skills: review-comet-pr, review-comet-shuffle-pr, and review-comet-memory-pr. Read the supplied discussions and checked current inline comments.

Exact-head CI: run 36732488695 was cancelled during preflight's Install mermaid-cli step. Required Checks failed, and build/Spark SQL jobs were skipped.

Validation: Compiled the exact-head rule and extracted current union implementation with the production union shim against Spark 4.1.3. The bounded planner/RDD harness passed, using minimal base/aggregate adapters and synthetic shuffle statistics. Full Comet/native builds, native query execution and multi-profile suites were not run. Reproduction: python3 /tmp/comet-6459-root-review/run.py. Project files remain unchanged.

case None => p
}
}
val coalesced = CoalesceShufflePartitions(SparkSession.active).apply(asSpark)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Preserve ancestor-join context when invoking Spark's coalescer. On Spark 4.0+, a CometUnionExec containing a coalescible shuffle and a scan can sit below CartesianProductExec, for example after a crossJoin with broadcasting disabled. Calling Spark's rule only on the union subtree hides the Cartesian ancestor, so Spark loses hasExplodingJoin and uses the ordinary advisory size instead of its smaller join-specific target. With 64 partitions reporting 1 MiB each, minPartitionNum=1, a 64 MiB advisory size and a 1 MiB minimum size, Spark and pre-PR Comet retain 64 partitions, while this rule reduces them to one. This concentrates the shuffled branch's cross-product work into one partition instead of 64. Please retain the containing join context when delegating, or skip these subtrees until that context can be preserved.

Evidence: Reproduced on Spark 4.1.3 with the exact-head rule and extracted current CometUnionExec. The plan is CartesianProductExec(CometUnionExec(ShuffleQueryStageExec, scan), scan) with synthetic statistics of 64 × 1 MiB. python3 /tmp/comet-6459-root-review/run.py exits 0 and reports EXPLODING_SPARK_COUNTS=Vector(64); EXPLODING_BEFORE_PR_COUNTS=Vector(64); EXPLODING_AFTER_PR_COUNTS=Vector(1). Source and output are in CurrentPlanProbe.scala and probe.log in that directory. Spark 4.0.4, 4.1.3 and 4.2.0 sources propagate hasExplodingJoin from ancestors and use the minimum partition size for these groups. This establishes a partition-parallelism regression without claiming a measured wall-clock slowdown.

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

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

AQE does not coalesce shuffle partitions under a CometUnion when a branch is not a shuffle

3 participants