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

   ## Which issue does this PR close?
   
   Part of #6133. This fixes the q5 failure. The q64 row loss has a separate 
cause, tracked in #6264.
   
   ## Rationale for this change
   
   AQE applies Spark's own query stage rules before Comet's. When 
`CoalesceShufflePartitions` changes a stage, `optimizeQueryStage` runs 
`ValidateRequirements`, which reads `outputPartitioning` on every node in the 
stage. `CometPlanAdaptiveDynamicPruningFilters` has not run yet at that point, 
so a scan's DPP subquery is still Spark's `SubqueryAdaptiveBroadcastExec` 
placeholder.
   
   For a non-bucketed scan, `CometNativeScanExec.outputPartitioning` returned 
`UnknownPartitioning(perPartitionData.length)`. Computing that length runs 
`InSubqueryExec.updateResult` on the placeholder, which throws 
`SubqueryAdaptiveBroadcastExec does not support the execute() code path`. That 
is the stack trace in the issue.
   
   Coalescing skips shuffle reads that share a subtree with a scan, and 
`CoalesceShufflePartitions` only splits a stage into separate groups at Spark's 
`UnionExec`, `BroadcastHashJoinExec`, `BroadcastNestedLoopJoinExec` and 
`CartesianProductExec`. So the failure needs one of those in the scan's stage. 
q5 has a union and hits it when Comet's union is not used.
   
   `FileSourceScanExec` does not have this problem: it reports 
`UnknownPartitioning(0)` for a non-bucketed scan and only resolves DPP when it 
builds its RDD.
   
   ## What changes are included in this PR?
   
   - `CometNativeScanExec.outputPartitioning` reports `UnknownPartitioning(0)` 
for a non-bucketed scan, as `FileSourceScanExec` does. A bucketed scan still 
reports its bucket layout.
   - `CometNativeExec` sizes the native RDD of a block whose leaf is a 
`CometScanWithPlanData` from `perPartitionData.length`, at execution, after 
`findAllPlanData` has resolved the scan's subqueries. It used to read that 
count from `outputPartitioning`. A bucketed scan has one entry per bucket, so 
its count is unchanged.
   
   Planning now sees 0 partitions for a non-bucketed native scan, the same as 
for Spark's scan. In particular a single split scan no longer satisfies 
`AllTuples` while planning, which also matches Spark. The execution paths that 
need a count take it from the RDD.
   
   ## How are these changes tested?
   
   Two tests in `CometExecSuite`, next to the existing AQE DPP tests:
   
   - "AQE DPP: shuffle coalescing in the scan's stage keeps DPP working" joins 
a partitioned fact table to a filtered dimension and unions it with an 
aggregate, with Comet's union disabled so Spark's `UnionExec` keeps the 
aggregate's shuffle read in the scan's stage. On main it fails with the stack 
trace from the issue. With the change it matches Spark, and on Spark 3.5 and 
later it checks that the DPP scan and a coalesced shuffle read share the final 
stage, and that DPP still prunes to 1 of 10 partitions.
   - "AQE DPP: bucketed scan keeps its partitioning for a sort-merge join" 
checks that a bucketed scan with DPP still feeds a sort-merge join without a 
shuffle. It passes on main and guards the bucketed branch.
   
   `CometExecSuite`, `CometScanWithPlanDataSuite`, `CometTopKSuite`, 
`CometJoinSuite` and the DPP fallback suites pass on Spark 3.5. The AQE DPP 
tests and `CometTopKSuite` pass on Spark 4.0 and 4.1.
   


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