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

   ### Describe the bug
   
   With AQE off, a dynamic partition pruning subquery doesn't reuse the join's 
`CometBroadcastExchange` when the dimension is a `UNION ALL`. The dimension is 
computed and broadcast a second time for the DPP subquery.
   
   `CometExecRule.rewriteInSubqueryPlan` only replaces the 
`SubqueryBroadcastExec` with a `CometSubqueryBroadcastExec` when the exchange's 
child, after stripping the columnar-to-row transition, is a `CometNativeExec`. 
`CometUnionExec` extends `CometExec` but not `CometNativeExec`, so the check 
fails and the subquery keeps a Spark `BroadcastExchange` over 
`CometColumnarToRow(CometUnion(...))`. `ReuseExchangeAndSubquery` can't match 
that to the join's `CometBroadcastExchange`. `CometCoalesceExec`, 
`CometTakeOrderedAndProjectExec` and a bare `CometSparkToColumnarExec` would 
hit the same check.
   
   The same query with AQE on reuses the broadcast through 
`CometPlanAdaptiveDynamicPruningFilters`. #6786 adds the same `CometNativeExec` 
check to that rule, which would make the AQE path behave like this one too.
   
   ### Steps to reproduce
   
   Reproduced on `main` at f9dd86d482 (Spark 4.1) with Comet's test session:
   
   ```scala
   spark.range(0, 1000).selectExpr("CAST(id % 10 AS STRING) AS fk", "id AS v")
     .write.partitionBy("fk").parquet(fact)
   spark.range(0, 5).selectExpr("CAST(id AS STRING) AS value", "id AS 
w").write.parquet(d1)
   spark.range(5, 10).selectExpr("CAST(id AS STRING) AS value", "id AS 
w").write.parquet(d2)
   
   // spark.sql.adaptive.enabled=false, 
spark.sql.optimizer.dynamicPartitionPruning.enabled=true
   sql("""SELECT /*+ BROADCAST(d) */ f.v FROM u_fact f JOIN
          (SELECT value FROM u_d1 WHERE w < 2 UNION ALL SELECT value FROM u_d2 
WHERE w > 8) d
          ON f.fk = d.value""")
   ```
   
   The executed plan:
   
   ```
   CometBroadcastHashJoin [cast(fk#34 as bigint)], [cast(value#35 as bigint)], 
Inner, BuildRight
   :- CometNativeScan parquet [v#33L,fk#34] ...
   :     +- SubqueryBroadcast dynamicpruning#39, [0], [cast(value#35 as bigint)]
   :        +- BroadcastExchange HashedRelationBroadcastMode(...)
   :           +- CometColumnarToRow
   :              +- CometUnion Union, [value#35]
   :                 :- CometProject / CometFilter / CometNativeScan parquet 
(d1)
   :                 +- CometProject / CometFilter / CometNativeScan parquet 
(d2)
   +- CometBroadcastExchange [value#35]
      +- CometUnion Union, [value#35]
         :- ... (same d1 branch)
         +- ... (same d2 branch)
   ```
   
   There is no `ReusedExchange` in the plan. With AQE on, the DPP subquery is a 
`CometSubqueryBroadcast` over `ReusedExchange [value#12], 
CometBroadcastExchange [value#12]`.
   
   ### Expected behavior
   
   The DPP subquery reuses the join's `CometBroadcastExchange`, as it does for 
a dimension whose top operator is a `CometNativeExec`, and as the AQE path does 
today. Results are correct either way. The cost is a second evaluation of the 
dimension and a second broadcast job.
   
   ### Additional context
   
   Found while reviewing #6786. One possible fix is to accept any Comet 
columnar plan in that check, for example `cometChild.isInstanceOf[CometPlan] && 
cometChild.supportsColumnar`. That change, applied to the AQE rule in #6786, 
restored reuse there and kept the existing DPP tests in `CometExecSuite` 
passing.
   


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