Skip to content

AQE skew join split makes the join fall back to Spark with no fallback reason #6530

Description

@comphead

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

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:shuffleShuffle (JVM and native)bugSomething isn't workingperformancepriority:mediumFunctional bugs, performance regressions, broken features

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions