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]