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

   ## Which issue does this PR close?
   
   Closes #6539.
   
   ## Rationale for this change
   
   Under AQE, the final `CometHashAggregate` of a two-phase aggregate shares 
its logical node with the shuffle stage below it, so 
`LogicalQueryStageStrategy` hands back the same physical node. Its native plan 
was serialized in the initial plan, when the input was a bare exchange and was 
written as a plain `Scan`. `CometExecRule` kept that plan, so the reduce side 
decoded every block on the JVM and imported it over FFI, even with 
`spark.comet.shuffle.directRead.enabled` on. Joins are planned fresh and 
already read natively.
   
   ## What changes are included in this PR?
   
   - `CometExecRule` refreshes a `CometNativeExec` whose plan still reads a 
child as a `Scan` when that child is now a `ShuffleScan` sink. The parent's 
leaves are matched to its children in the order the native block reads its 
inputs, a leaf is replaced only when its field types equal the `ShuffleScan`'s, 
and any mismatch leaves the node as it was. Nested native children propagate 
their leaves to the block root. The three write operators, which build their 
own `Scan` child, are skipped. The leaf is patched in place rather than 
converting again from `originalPlan`, which would drop the stage's logical link 
and run serde again on a planned node.
   - `CometNativeExec.withRefreshedNativeOp` copies the node with the new plan 
and an empty serialized plan, so the block is serialized again.
   - `native_shuffle.md` notes the second point where direct read is chosen.
   
   This also covers reads that AQE coalesces: the aggregate takes the stage's 
`ShuffleScan` before AQE adds the coalesced `AQEShuffleReadExec`, and the input 
wiring already accepts that read. The skew-join fallback in #6530 is a separate 
path and is not changed here.
   
   Reduce-stage median and end-to-end time from a local benchmark (container 
with 8 CPUs and 7 GB, `local[4]`, 5 timed runs), `main` at 351facd26 against 
this PR. Results matched Spark in every case.
   
   | Case | Spark | Comet `main` | Comet with this PR |
   |---|---|---|---|
   | 2000 map tasks, 200 shuffle partitions, AQE coalescing off | 1,441 / 7,535 
ms | 2,783 / 5,172 ms | 1,291 / 2,976 ms |
   | same, AQE coalescing on | 1,635 / 7,589 ms | 3,263 / 5,153 ms | 1,496 / 
3,175 ms |
   | 2000 map tasks, 2000 shuffle partitions, AQE coalescing on | 10,682 / 
15,778 ms | 24,831 / 29,853 ms | 10,364 / 15,612 ms |
   | 16 map tasks, 200 shuffle partitions | 96 / 362 ms | 189 / 344 ms | 194 / 
359 ms |
   
   The query is `spark.range(0, 20000000, 1, maps).groupBy(col("id") % 
50000).agg(sum("id"), count("id"))`.
   
   ## How are these changes tested?
   
   New tests in `CometNativeShuffleSuite` and `CometExecRuleSuite` check which 
leaf each native block reads, paired with its input in order, and compare 
answers with Spark:
   
   - a final aggregate over a shuffle stage reads `ShuffleScan`, with AQE 
coalescing off and on;
   - a chain of native aggregates, and a join of final aggregates where one 
side reuses the exchange and another variant adds a broadcast side (only the 
shuffle slot changes);
   - an AQE DPP query with a broadcast side (Spark 3.5+);
   - exchange reuse still happens above two equivalent refreshed aggregates;
   - direct read off and AQE off keep a plain `Scan`;
   - at the rule level, a second pass changes nothing, and a field-type 
mismatch leaves the node unchanged.
   
   The #5483 test "AQE DPP broadcast roots retain temporary logical links after 
an unchanged replan" now runs with direct read off, since with it on the first 
replan is no longer unchanged. A sibling test covers the default: the replan 
that moves the aggregate under the broadcast onto a `ShuffleScan` is adopted, 
and every replanned broadcast root keeps its logical link and 
`TEMP_LOGICAL_PLAN_TAG`.
   
   The new tests fail on `main`.
   


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