Skip to content

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

Description

@andygrove

Describe the bug

With AQE on, a union that Comet converts to CometUnion gets 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. The shuffled branch keeps spark.sql.shuffle.partitions partitions, so a small union runs hundreds of near-empty tasks.

Spark's CoalesceShufflePartitions treats each child of a UnionExec, CartesianProductExec, BroadcastHashJoinExec or BroadcastNestedLoopJoinExec as an independent coalesce group (childrenNeedCompatiblePartitioning). It matches those classes directly, so a CometUnionExec falls through to the case that coalesces only when every leaf under it is an exchange stage. When a branch ends in a scan, nothing is coalesced. Comet converts the union before the rule runs, because CometRule(session, queryStagePrep = true) is registered as a query-stage preparation rule, and CoalesceShufflePartitions is a query-stage optimizer rule that runs later.

I found this through Spark's SPARK-42101: Coalesce shuffle partition with union even if exists TableCacheQueryStage while running Spark's SQL suites with Comet's cache format, but it does not depend on the cache.

Steps to reproduce

On main at 8e4bded, with AQE on and spark.sql.adaptive.coalescePartitions.minPartitionNum=1:

spark.range(0, 100, 1, 1).toDF("c").write.parquet(path)
val df = spark.range(0, 10, 1, 2).toDF("c").repartition($"c").union(spark.read.parquet(path))
df.collect()
df.rdd.getNumPartitions

With Comet the final plan is CometUnion over an uncoalesced ShuffleQueryStage and a CometNativeScan, and the result has 201 partitions. With Comet disabled the shuffle branch is read through AQEShuffleRead coalesced, and the result has 2.

Expected behavior

The shuffle branches under a CometUnion are coalesced as they are under Spark's UnionExec.

Additional context

The same class match probably affects CometBroadcastHashJoinExec and any other Comet counterpart of the four operators above, though I have not checked them. The fix could defer the union conversion until after CoalesceShufflePartitions has run, or have Comet run the equivalent coalescing for its own operators as a query-stage optimizer rule.

Activity

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

Metadata

Metadata

Assignees

Labels

area:shuffleShuffle (JVM and native)bugSomething isn't workingperformancepriority:mediumFunctional bugs, performance regressions, broken features

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions