Describe the bug
With AQE on, the final CometHashAggregate of a two-phase aggregate never uses the native shuffle direct read (spark.comet.shuffle.directRead.enabled, on by default). Every shuffle block is decoded on the JVM and imported over FFI instead.
The final aggregate shares its logical node with the shuffle stage below it, so AQE's replaceWithQueryStagesInLogicalPlan wraps it in a LogicalQueryStage, and LogicalQueryStageStrategy hands back the same physical node. CometExecRule then keeps that node's native plan, which was serialized in the initial plan, when its input was still a bare exchange and was written as a plain Scan. CometExchangeSink.shouldUseShuffleScan returns true for the stage, but the stale plan never picks it up. This happens with AQE partition coalescing on or off. Joins have their own logical node, are planned fresh, and do use direct read.
Steps to reproduce
spark.range(0, 20000000, 1, 2000)
.groupBy(col("id") % 50000)
.agg(sum("id"), count("id"))
.collect()
With AQE on, the final aggregate's native plan is Projection -> HashAgg -> Scan, and it opens no native shuffle scans.
Expected behavior
The final aggregate reads its shuffle through ShuffleScanExec, as joins do.
Additional context
Refreshing the stale leaf in a prototype, measured locally (container with 8 CPUs, local[4], reduce-stage median):
| Case |
Spark |
Comet |
Comet with refresh |
| 2000 map tasks, 200 shuffle partitions |
1,316 ms |
2,746 ms |
1,175 ms |
| same, AQE partition coalescing on |
1,601 ms |
3,394 ms |
1,421 ms |
| 2000 map tasks, 2000 shuffle partitions, AQE coalescing on |
10,434 ms |
24,849 ms |
10,223 ms |
Results matched Spark in every case, including an AQE skew-split join and a runtime broadcast join. The skew-join fallback itself is #6530. At 2000 partitions the map side is the remaining cost, which #5905 covers.
Describe the bug
With AQE on, the final
CometHashAggregateof a two-phase aggregate never uses the native shuffle direct read (spark.comet.shuffle.directRead.enabled, on by default). Every shuffle block is decoded on the JVM and imported over FFI instead.The final aggregate shares its logical node with the shuffle stage below it, so AQE's
replaceWithQueryStagesInLogicalPlanwraps it in aLogicalQueryStage, andLogicalQueryStageStrategyhands back the same physical node.CometExecRulethen keeps that node's native plan, which was serialized in the initial plan, when its input was still a bare exchange and was written as a plainScan.CometExchangeSink.shouldUseShuffleScanreturns true for the stage, but the stale plan never picks it up. This happens with AQE partition coalescing on or off. Joins have their own logical node, are planned fresh, and do use direct read.Steps to reproduce
With AQE on, the final aggregate's native plan is
Projection -> HashAgg -> Scan, and it opens no native shuffle scans.Expected behavior
The final aggregate reads its shuffle through
ShuffleScanExec, as joins do.Additional context
Refreshing the stale leaf in a prototype, measured locally (container with 8 CPUs,
local[4], reduce-stage median):Results matched Spark in every case, including an AQE skew-split join and a runtime broadcast join. The skew-join fallback itself is #6530. At 2000 partitions the map side is the remaining cost, which #5905 covers.