andygrove opened a new pull request, #6459:
URL: https://github.com/apache/datafusion-comet/pull/6459

   ## Which issue does this PR close?
   
   Closes #6454.
   
   ## Rationale for this change
   
   With AQE on, a union that Comet converts to `CometUnionExec` got 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, so the shuffled branches kept 
`spark.sql.shuffle.partitions` partitions. The issue has the analysis and a 
reproduction that ends with 201 partitions where Spark has 2.
   
   Spark's `CoalesceShufflePartitions` coalesces each child of a `UnionExec` as 
its own group, and from Spark 4.0 each child of a `CartesianProductExec`, 
`BroadcastHashJoinExec` or `BroadcastNestedLoopJoinExec` too, but it matches 
those classes. Comet replaces them while AQE prepares a stage, before the 
optimizer rules run, so Spark's rule sees the Comet operators and falls through 
to the case that coalesces only when every leaf below is an exchange stage.
   
   ## What changes are included in this PR?
   
   A new query-stage optimizer rule, `CometCoalesceShufflePartitions`, 
registered after Comet's other two, so it runs after Spark's 
`CoalesceShufflePartitions`. It looks for Comet operators whose `originalPlan` 
is one of the operators above, and whose shuffle stages no AQE rule has put a 
read over. For each one it rebuilds the Spark operators over the Comet 
children, runs Spark's own `CoalesceShufflePartitions` over that, and swaps the 
Comet operators back in over the new children. Spark's code makes every 
coalescing decision, per Spark version.
   
   The rule extends `AQEShuffleReadRule`, so AQE treats it like Spark's rule: 
it is skipped for the final stage when 
`spark.sql.adaptive.applyFinalStageShuffleOptimizations` is off, and its result 
is discarded if it breaks a distribution required above it.
   
   It also covers a union whose branches are all shuffles with different 
partition counts. Spark's rule coalesces a Comet union's shuffles together, 
finds the counts differ, and coalesces nothing, where for a Spark union it 
would coalesce each branch on its own.
   
   Two limits:
   
   - When every leaf below a Comet union is an exchange stage and the counts 
agree, Spark's rule already coalesces its shuffles, together rather than branch 
by branch, and this rule leaves them as they are.
   - Spark 3.4 has no hook for query-stage optimizer rules, so nothing changes 
there.
   
   Once this lands, #5634 can drop the `IgnoreComet` it puts on Spark's 
`SPARK-42101` union test.
   
   ## How are these changes tested?
   
   - A new test in `CometExecSuite` runs the issue's query, a union of a 
repartitioned DataFrame with a Parquet scan. It checks for a coalesced read 
under the `CometUnionExec`, and for the same number of result partitions as 
with Comet disabled.
   - A port of Spark's `SPARK-42101: Coalesce shuffle partition with union even 
if exists TableCacheQueryStage` to `CometInMemoryCacheSuite`.
   
   Both fail with the rule unregistered. They pass on Spark 3.5, 4.0, 4.1 and 
4.2, and are canceled on 3.4.
   
   Spark 4.1's `AdaptiveQueryExecSuite`, `CoalesceShufflePartitionsSuite`, 
`DataFrameSetOperationsSuite`, `ExchangeSuite` and `DataFrameSuite`, run from 
the test jar with Comet enabled, give the same results with and without the 
rule, except that `CoalesceShufflePartitionsSuite`'s `Union two datasets with 
different pre-shuffle partition number` fails without it and passes with it.
   
   `CometExecSuite`, `CometJoinSuite`, `CometSetOpWithGroupBySuite`, 
`CometPlanEqualitySuite`, `CometInMemoryCacheSuite` and 
`CometNativeShuffleSuite` pass on 4.1. The strict-warnings compile and scalafix 
pass on 3.5.
   


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