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

   ### What is the problem the feature request solves?
   
   Without AQE, a native operator that reads a Comet shuffle gets every block 
through `NativeBatchDecoderIterator`. Native code decodes the block, the JVM 
imports it into vectors, and the JVM exports it back to the native plan through 
the Arrow C stream. Behind an AQE query stage, the same operator reads the 
block directly in native code instead (`ShuffleScan`, through 
`CometShuffleBlockIterator`). The JVM leg costs a fixed amount per block, which 
dominates reduce tasks that read many small blocks with a wide schema.
   
   The test case is a shuffled hash join over rows with 3 nested and 100 flat 
columns, with 32 map tasks, 32 reduce partitions and 95% of the rows on one 
key. Each of the 31 small reduce tasks reads 32 blocks of about 100 rows:
   
   | Small reduce task | Spark | Comet |
   |---|---:|---:|
   | AQE off, mean in the join stage | 16.5 ms | 65 ms |
   | AQE off, the same task run alone | 19 ms | 30 ms |
   | AQE on (direct read), mean in the join stage | 17–26 ms | 13–14 ms |
   
   In other runs without AQE, the stage means were 17–33 ms for Spark and 60–64 
ms for Comet.
   
   Without AQE the plan reads none of its 2 shuffle inputs directly, and with 
AQE it reads both (`CometExec.findShuffleScanIndices` on the block's native 
plan). The comment at 
[`operators.scala:1040-1044`](https://github.com/apache/datafusion-comet/blob/9a6241f140e9fd576ba9467ee26c2e21cf46ec02/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala#L1040-L1044)
 gives the reason: `CometExchangeSink.shouldUseShuffleScan` only fires for AQE 
wrappers (`ShuffleQueryStageExec`), so a bare non-AQE 
`CometShuffleExchangeExec` always serializes as a regular Scan, whatever 
`spark.comet.shuffle.directRead.enabled` says.
   
   Profile of one small Comet task, run alone about 800 times, from a 1 ms 
stack sampler on the task thread plus macOS `sample`:
   
   | Share | Where the Comet task's time goes |
   |---:|---|
   | 47% | The JVM leg: import into JVM vectors 13%, Comet vector objects 9%, 
export back to native 14%, `org.apache.arrow.vector.ipc` message objects 6%, 
JNI releases 5% |
   | 22% | Native plan: join, take, probe, import |
   | 12% | Native decode (`decodeShuffleBlock`) |
   | 11% | Spark's fetch path, which opens and reads the shuffle file for each 
block |
   | 8% | Other |
   
   - Arrow's JVM buffer allocation and release (`Unsafe.allocateMemory0`, 
`Unsafe.freeMemory0`, `AllocationManager`) takes about 22% of the task by 
itself, because every block carries hundreds of buffers for its 104 columns.
   - macOS `sample` attributes 11.6% of the task to the W^X toggle that macOS 
on Apple Silicon makes at every JNI transition. That part would not exist on 
Linux. Another 6% is in `System.nanoTime`.
   
   Spark's small task is about a third reading and copying rows, 20% lz4-java 
decompression and 19% output row writes.
   
   Comet's small tasks also take twice as long in the join stage as alone (65 
against 30 ms), while Spark's do not. I have not investigated why.
   
   ### Describe the potential solution
   
   - Let native consumers read a Comet shuffle directly in native code without 
AQE too, as they already do behind `ShuffleQueryStageExec`. That removes the 
JVM leg for every block, and with AQE these same tasks already beat Spark.
   - Short of that, make the JVM leg cheaper per block. Most of it allocates 
and releases Arrow buffers and builds vector objects for a batch that no JVM 
code reads.
   
   ### Additional context
   
   - Related to #5905. Its R3 covers the two FFI round trips per batch when a 
JVM operator consumes the shuffle, and its R5 the per-block buffer reallocation 
in `NativeBatchDecoderIterator`. This issue is the native-consumer case without 
AQE, where no JVM code needs the batch at all.
   - Also related: #6539 (closed), the same JVM decoder path for a final 
aggregate under AQE.
   - Found while profiling #6528.
   - To reproduce:
     - `fact`: 2,097,152 rows in 32 Parquet files, with a struct holding an 
array of structs, an array of strings, a map and 100 flat columns.
     - `dim`: 262,144 rows.
     - Query: `SELECT /*+ SHUFFLE_HASH(d) */ ... FROM fact f JOIN dim d ON 
f.key = d.k`.
     - Settings: `local[8]`, `spark.sql.shuffle.partitions=32`, AQE off. 
Comet's output is counted from the native plan's batches, without converting it 
to rows.
     - Measured with a local benchmark that is not in a PR 
(`CometHashJoinTaskTimeBenchmark -- profile=small`).
   - Environment: Comet `main` at 9a6241f140 with the unmerged #6805, which 
changes only how the shuffle readers read their stream. Spark 4.1.3, Scala 
2.13, JDK 17.0.19, macOS 15.7.4 on an Apple M3 Max, release build.
   
   ### 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