Conversation
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.
comphead
left a comment
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
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) { |
There was a problem hiding this comment.
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") { |
There was a problem hiding this comment.
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 | |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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
AQEShuffleReadRulepreserves 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) |
There was a problem hiding this comment.
[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.
Which issue does this PR close?
Closes #6454.
Rationale for this change
With AQE on, a union that Comet converts to
CometUnionExecgot 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 keptspark.sql.shuffle.partitionspartitions. The issue has the analysis and a reproduction that ends with 201 partitions where Spark has 2.Spark's
CoalesceShufflePartitionscoalesces each child of aUnionExecas its own group, and from Spark 4.0 each child of aCartesianProductExec,BroadcastHashJoinExecorBroadcastNestedLoopJoinExectoo, 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'sCoalesceShufflePartitions. It looks for Comet operators whoseoriginalPlanis 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 ownCoalesceShufflePartitionsover 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 whenspark.sql.adaptive.applyFinalStageShuffleOptimizationsis 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:
Once this lands, #5634 can drop the
IgnoreCometit puts on Spark'sSPARK-42101union test.How are these changes tested?
CometExecSuiteruns the issue's query, a union of a repartitioned DataFrame with a Parquet scan. It checks for a coalesced read under theCometUnionExec, and for the same number of result partitions as with Comet disabled.SPARK-42101: Coalesce shuffle partition with union even if exists TableCacheQueryStagetoCometInMemoryCacheSuite.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,ExchangeSuiteandDataFrameSuite, run from the test jar with Comet enabled, give the same results with and without the rule, except thatCoalesceShufflePartitionsSuite'sUnion two datasets with different pre-shuffle partition numberfails without it and passes with it.CometExecSuite,CometJoinSuite,CometSetOpWithGroupBySuite,CometPlanEqualitySuite,CometInMemoryCacheSuiteandCometNativeShuffleSuitepass on 4.1. The strict-warnings compile and scalafix pass on 3.5.