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

   ### Describe the bug
   
   With AQE on, a native operator above a Comet shuffle reads that shuffle 
directly in native code through `ShuffleScanExec` 
(`spark.comet.shuffle.directRead.enabled`, on by default). None of the scan's 
timing reaches Spark's SQL metrics.
   
   In a reduce task that read 1,512 MiB of shuffle data in 1,411 ms, the 
`CometExchange` nodes reported `local bytes read`, `records read` and `fetch 
wait time`, but no updates at all for `native shuffle writer time` 
(`elapsed_compute`) or `decoding and decompression time` (`decode_time`). No 
other node reported them either. `CometHashJoin` reported 404 ms of 
`join_time`. Most of the remaining second, reading shuffle blocks, lz4 
decompression and IPC decode, cannot be attributed from the operator metrics.
   
   The JVM side explains it:
   
   - `CometNativeExec` builds the task's metric tree with 
[`CometMetricNode.fromCometPlan(this)`](https://github.com/apache/datafusion-comet/blob/7d294535e87367d819c4e2a23422cd1158ac31b8/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala#L930),
 which [follows the Spark plan's 
`children`](https://github.com/apache/datafusion-comet/blob/7d294535e87367d819c4e2a23422cd1158ac31b8/spark/src/main/scala/org/apache/spark/sql/comet/CometMetricNode.scala#L580).
   - Under AQE, the shuffle input of a native block is a 
`ShuffleQueryStageExec`. In Spark, `QueryStageExec extends LeafExecNode` and 
defines no metrics (checked in 3.4.3, 3.5.8, 4.0.1 and 4.1.3), so its node in 
the tree has an empty metric map.
   - 
[`set_all`](https://github.com/apache/datafusion-comet/blob/7d294535e87367d819c4e2a23422cd1158ac31b8/spark/src/main/scala/org/apache/spark/sql/comet/CometMetricNode.scala#L246)
 matches native children to these nodes by position, and 
[`set`](https://github.com/apache/datafusion-comet/blob/7d294535e87367d819c4e2a23422cd1158ac31b8/spark/src/main/scala/org/apache/spark/sql/comet/CometMetricNode.scala#L232)
 ignores a name that is not in the node's map. So every metric 
`ShuffleScanExec` reports (`elapsed_compute`, `decode_time`, `output_rows`) is 
dropped.
   - Without AQE, the input is the `CometShuffleExchangeExec` itself, whose map 
has these names, so the time is recorded there.
   
   I have not traced how the native side lays out the metric tree, so the above 
is inferred from the JVM code and the observed metrics.
   
   Two related symptoms, likely in scope of the same fix:
   
   - Without AQE, the reduce-side read time is recorded under the writer's 
display name, `native shuffle writer time`, because the reduce-side `ScanExec` 
reports `elapsed_compute` into the exchange's writer metric.
   - With partition coalescing on (the default), the shuffle input is probably 
wrapped in an `AQEShuffleReadExec`, whose metrics do not have these names 
either. Not tested.
   
   ### Steps to reproduce
   
   1. Join two Parquet tables with a shuffled hash join whose output a native 
operator consumes, for example `SELECT /*+ SHUFFLE_HASH(d) */ ... FROM fact f 
JOIN dim d ON f.key = d.k`.
   2. Run with Comet and Comet shuffle enabled, 
`spark.sql.adaptive.enabled=true`, `spark.sql.adaptive.skewJoin.enabled=false` 
and `spark.sql.adaptive.coalescePartitions.enabled=false`. The last two only 
keep the plan simple.
   3. Read the reduce tasks' SQL metric updates, for example from 
`taskInfo.accumulables` in a `SparkListener` or from the event log. The native 
block reads both inputs through `ShuffleScan` 
(`CometExec.findShuffleScanIndices` returns 2 for its native plan), yet the 
reduce tasks carry no `elapsed_compute` or `decode_time` updates for either 
exchange.
   
   Observed with a local benchmark (`CometHashJoinTaskTimeBenchmark -- aqe`, 
not in a PR yet): 2,097,152 fact rows with 3 nested and 100 flat columns, 95% 
of them on one key, and 32 shuffle partitions.
   
   ### Expected behavior
   
   The time `ShuffleScanExec` spends reading and decoding shuffle blocks 
appears in the SQL metrics of the exchange it reads (or of its query stage), 
with and without AQE, under names that say read rather than write.
   
   ### Additional context
   
   - Found while profiling #6528, which needs these metrics to attribute the 
time of its hot task.
   - A possible direction: in `CometMetricNode.fromCometPlan`, give a 
`ShuffleQueryStageExec` (including one over a `ReusedExchangeExec`) and an 
`AQEShuffleReadExec` the metrics of the `CometShuffleExchangeExec` they wrap, 
and report reduce-side read time under its own name instead of `native shuffle 
writer time`.
   - Environment: Comet `main` at 7d294535e8, Spark 4.1.3, Scala 2.13, JDK 
17.0.19, macOS 15.7.4 on an Apple M3 Max.
   
   ### Willingness to contribute
   
   I can contribute a fix for this bug independently
   


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