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.
Describe the bug
With AQE on, a union that Comet converts to
CometUniongets 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 keepsspark.sql.shuffle.partitionspartitions, so a small union runs hundreds of near-empty tasks.Spark's
CoalesceShufflePartitionstreats each child of aUnionExec,CartesianProductExec,BroadcastHashJoinExecorBroadcastNestedLoopJoinExecas an independent coalesce group (childrenNeedCompatiblePartitioning). It matches those classes directly, so aCometUnionExecfalls 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, becauseCometRule(session, queryStagePrep = true)is registered as a query-stage preparation rule, andCoalesceShufflePartitionsis a query-stage optimizer rule that runs later.I found this through Spark's
SPARK-42101: Coalesce shuffle partition with union even if exists TableCacheQueryStagewhile running Spark's SQL suites with Comet's cache format, but it does not depend on the cache.Steps to reproduce
On
mainat 8e4bded, with AQE on andspark.sql.adaptive.coalescePartitions.minPartitionNum=1:With Comet the final plan is
CometUnionover an uncoalescedShuffleQueryStageand aCometNativeScan, and the result has 201 partitions. With Comet disabled the shuffle branch is read throughAQEShuffleRead coalesced, and the result has 2.Expected behavior
The shuffle branches under a
CometUnionare coalesced as they are under Spark'sUnionExec.Additional context
The same class match probably affects
CometBroadcastHashJoinExecand any other Comet counterpart of the four operators above, though I have not checked them. The fix could defer the union conversion until afterCoalesceShufflePartitionshas run, or have Comet run the equivalent coalescing for its own operators as a query-stage optimizer rule.