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

   ### 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 d5de7b9927)
   
   - 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](https://github.com/apache/datafusion-comet/blob/d5de7b9927996cb7a87d84e658dade523e197c49/spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala#L518)),
 but has no case for an `AQEShuffleReadExec` wrapping that stage. The read is 
left as it is 
([L549](https://github.com/apache/datafusion-comet/blob/d5de7b9927996cb7a87d84e658dade523e197c49/spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala#L549)).
   - The default branch only tries an operator's serde when all of its children 
are `CometNativeExec` 
([L537](https://github.com/apache/datafusion-comet/blob/d5de7b9927996cb7a87d84e658dade523e197c49/spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala#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](https://github.com/apache/datafusion-comet/blob/d5de7b9927996cb7a87d84e658dade523e197c49/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffledRowRDD.scala#L95),
 
[L124](https://github.com/apache/datafusion-comet/blob/d5de7b9927996cb7a87d84e658dade523e197c49/spark/src/main/scala/org/apache/spark/sql/comet/execution/shuffle/CometShuffledRowRDD.scala#L124)),
 and `RewriteJoin` keeps `isSkewJoin` 
([L86](https://github.com/apache/datafusion-comet/blob/d5de7b9927996cb7a87d84e658dade523e197c49/spark/src/main/scala/org/apache/comet/rules/RewriteJoin.scala#L86)).
 So the gap looks to be plan conversion only. Note that 
`CometExchangeSink.shouldUseShuffleScan` 
([CometSink.scala#L108](https://github.com/apache/datafusion-comet/blob/d5de7b9927996cb7a87d84e658dade523e197c49/spark/src/main/scala/org/apache/comet/serde/operator
 /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]`:
   
   ```scala
   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 d5de7b9927 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
   
   - The "Comet AQE on" hot-key numbers for the shuffled hash join in #6528 
come from this fallback.
   - #6442 is a different case: there a join that Comet declined loses its 
reason. Here the join is never attempted.
   


-- 
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