Describe the bug
With AQE on, when Spark's OptimizeSkewedJoin splits a skewed partition of a sort-merge or shuffled hash join, Comet runs the reduce side of the join in Spark. The join becomes Spark's SortMergeJoin(skew=true) (or ShuffledHashJoin(skew=true)) with a ColumnarToRow over each AQEShuffleRead. Everything above the join in that stage runs in Spark too, and an aggregate on top cannot go native because its partial side is now a Spark aggregate.
Comet records no fallback reason for the join, so neither explain nor the fallback log says why. The scans and the shuffle writes stay native, and results are correct.
Plan from the repro below. With AQE off, all 13 operators are native:
CometColumnarToRow
+- CometHashAggregate
+- CometExchange
+- CometHashAggregate
+- CometProject
+- CometSortMergeJoin
:- CometSort
: +- CometExchange
: +- CometFilter
: +- CometNativeScan parquet
+- CometSort
+- CometExchange
+- CometFilter
+- CometNativeScan parquet
With AQE on, the hot partition is split and 6 of 13 operators are native:
HashAggregate [COMET: Comet aggregate that merges intermediate buffers requires a Comet child aggregate when the intermediate buffer formats are incompatible with Spark. Incompatible aggregate function(s): count]
+- Exchange
+- HashAggregate
+- Project
+- SortMergeJoin(skew=true)
:- Sort
: +- ColumnarToRow
: +- AQEShuffleRead
: +- CometExchange
: +- CometFilter
: +- CometNativeScan parquet
+- Sort
+- ColumnarToRow
+- AQEShuffleRead
+- CometExchange
+- CometFilter
+- CometNativeScan parquet
Cause (from reading the code at d5de7b9)
- Spark runs
OptimizeSkewedJoin in AdaptiveSparkPlanExec.queryStagePreparationRules and appends the extensions' query stage prep rules after it, in both 3.4 and 4.1. So when CometRule runs on the re-optimized plan, the join's children are already SortExec(AQEShuffleReadExec(ShuffleQueryStageExec(CometShuffleExchangeExec))).
CometExecRule.convertNode turns ShuffleQueryStageExec(_, _: CometShuffleExchangeExec, _) into a CometExchangeSink (CometExecRule.scala#L518), but has no case for an AQEShuffleReadExec wrapping that stage. The read is left as it is (L549).
- The default branch only tries an operator's serde when all of its children are
CometNativeExec (L537). The SortExec above the read is therefore never attempted, the join then has non-native children, and nothing gets a fallback reason.
- Without a skew split, AQE's
CoalesceShufflePartitions adds its AQEShuffleReadExec in queryStageOptimizerRules, after Comet has converted the plan, so the read lands under an existing CometExchangeSink. That is why plain AQE coalescing stays native.
- The read path already understands skew splits.
CometShuffledRowRDD handles PartialReducerPartitionSpec and PartialMapperPartitionSpec (L95, L124), and RewriteJoin keeps isSkewJoin (L86). So the gap looks to be plan conversion only. Note that CometExchangeSink.shouldUseShuffleScan (CometSink.scala#L108) also recognizes only a stage or the exchange itself, so a fix has to decide whether the direct read handles the partition specs or the sink falls back to the RDD read for an AQEShuffleReadExec.
Impact
From the skewed-join benchmark in #6528. Total task time from the event log, the median of 5 runs, in ms:
| case |
Spark AQE on |
Spark AQE off |
Comet AQE on |
Comet AQE off |
| SMJ then aggregate, flat, uniform keys |
2717 |
2676 |
875 |
1004 |
| SMJ then aggregate, flat, hot key 50% |
2825 |
2881 |
1859 |
995 |
| SMJ then aggregate, flat, hot key 90% |
2221 |
2445 |
1394 |
1010 |
| SMJ then aggregate, nested, hot key 50% |
4728 |
4271 |
3101 |
1655 |
| SMJ then aggregate, nested, hot key 90% |
4658 |
4519 |
2628 |
1633 |
| SMJ, nested, hot key 50% |
4500 |
4177 |
2573 |
2213 |
| SMJ, skewed side buffered, nested, hot key 90% |
4534 |
4688 |
2462 |
2039 |
With a hot key, turning AQE on costs Comet 38% to 87% more task time in the aggregation cases, while Spark's changes by 11% at most. Comet's lead over Spark there drops from 2.4x to 2.9x with AQE off, to 1.5x to 1.8x with AQE on. With uniform keys nothing is split and AQE on is as fast or faster for Comet. Plain joins lose less (up to 22%), because only the join and the projection above it move to Spark.
Steps to reproduce
Session with Comet enabled, native shuffle (spark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager) and off-heap memory, local[8]:
spark.range(0, 2097152, 1, 32)
.selectExpr("CAST(IF(id % 100 < 90, 0, 1 + PMOD(HASH(id), 262143)) AS BIGINT) AS k", "id AS v")
.write.parquet("/tmp/skew/fact") // 90% of rows on key 0
spark.range(0, 262144, 1, 8).selectExpr("id AS k", "id * 3 AS w").write.parquet("/tmp/skew/dim")
spark.read.parquet("/tmp/skew/fact").createOrReplaceTempView("fact")
spark.read.parquet("/tmp/skew/dim").createOrReplaceTempView("dim")
spark.conf.set("spark.comet.enabled", "true")
spark.conf.set("spark.comet.exec.enabled", "true")
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
// Scaled down so that a partition of this size counts as skewed
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "4MB")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "1MB")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "-1")
spark.conf.set("spark.sql.shuffle.partitions", "32")
val df = spark.sql(
"SELECT /*+ MERGE(d) */ COUNT(*), SUM(f.v + d.w) FROM fact f JOIN dim d ON f.k = d.k")
df.collect()
println(new org.apache.comet.ExtendedExplainInfo().generateVerboseInfo(df.queryExecution.executedPlan))
I ran this as a small program on main at d5de7b9 with Spark 4.1.3. Setting spark.sql.adaptive.enabled=false gives the fully native plan above. In the benchmark, the shuffled hash join (/*+ SHUFFLE_HASH(d) */) falls back the same way, as ShuffledHashJoin(skew=true).
Expected behavior
The skew-split join and the operators above it run natively, with the Comet join reading the split partitions through a native shuffle input, as it already does for coalesced AQE reads. If a skew split cannot be supported, the join should carry a fallback reason.
Additional context
Describe the bug
With AQE on, when Spark's
OptimizeSkewedJoinsplits a skewed partition of a sort-merge or shuffled hash join, Comet runs the reduce side of the join in Spark. The join becomes Spark'sSortMergeJoin(skew=true)(orShuffledHashJoin(skew=true)) with aColumnarToRowover eachAQEShuffleRead. Everything above the join in that stage runs in Spark too, and an aggregate on top cannot go native because its partial side is now a Spark aggregate.Comet records no fallback reason for the join, so neither
explainnor the fallback log says why. The scans and the shuffle writes stay native, and results are correct.Plan from the repro below. With AQE off, all 13 operators are native:
With AQE on, the hot partition is split and 6 of 13 operators are native:
Cause (from reading the code at d5de7b9)
OptimizeSkewedJoininAdaptiveSparkPlanExec.queryStagePreparationRulesand appends the extensions' query stage prep rules after it, in both 3.4 and 4.1. So whenCometRuleruns on the re-optimized plan, the join's children are alreadySortExec(AQEShuffleReadExec(ShuffleQueryStageExec(CometShuffleExchangeExec))).CometExecRule.convertNodeturnsShuffleQueryStageExec(_, _: CometShuffleExchangeExec, _)into aCometExchangeSink(CometExecRule.scala#L518), but has no case for anAQEShuffleReadExecwrapping that stage. The read is left as it is (L549).CometNativeExec(L537). TheSortExecabove the read is therefore never attempted, the join then has non-native children, and nothing gets a fallback reason.CoalesceShufflePartitionsadds itsAQEShuffleReadExecinqueryStageOptimizerRules, after Comet has converted the plan, so the read lands under an existingCometExchangeSink. That is why plain AQE coalescing stays native.CometShuffledRowRDDhandlesPartialReducerPartitionSpecandPartialMapperPartitionSpec(L95, L124), andRewriteJoinkeepsisSkewJoin(L86). So the gap looks to be plan conversion only. Note thatCometExchangeSink.shouldUseShuffleScan(CometSink.scala#L108) also recognizes only a stage or the exchange itself, so a fix has to decide whether the direct read handles the partition specs or the sink falls back to the RDD read for anAQEShuffleReadExec.Impact
From the skewed-join benchmark in #6528. Total task time from the event log, the median of 5 runs, in ms:
With a hot key, turning AQE on costs Comet 38% to 87% more task time in the aggregation cases, while Spark's changes by 11% at most. Comet's lead over Spark there drops from 2.4x to 2.9x with AQE off, to 1.5x to 1.8x with AQE on. With uniform keys nothing is split and AQE on is as fast or faster for Comet. Plain joins lose less (up to 22%), because only the join and the projection above it move to Spark.
Steps to reproduce
Session with Comet enabled, native shuffle (
spark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager) and off-heap memory,local[8]:I ran this as a small program on
mainat d5de7b9 with Spark 4.1.3. Settingspark.sql.adaptive.enabled=falsegives the fully native plan above. In the benchmark, the shuffled hash join (/*+ SHUFFLE_HASH(d) */) falls back the same way, asShuffledHashJoin(skew=true).Expected behavior
The skew-split join and the operators above it run natively, with the Comet join reading the split partitions through a native shuffle input, as it already does for coalesced AQE reads. If a skew split cannot be supported, the join should carry a fallback reason.
Additional context