andygrove opened a new issue, #6454:
URL: https://github.com/apache/datafusion-comet/issues/6454

   ### 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 8e4bded55, with AQE on and 
`spark.sql.adaptive.coalescePartitions.minPartitionNum=1`:
   
   ```scala
   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.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to