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]
