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]