You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
perf: without AQE, native operators read Comet shuffle through a JVM round trip per block, making small reduce tasks up to 4x slower than Spark #6807
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 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.
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, throughCometShuffleBlockIterator). 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:
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.findShuffleScanIndiceson the block's native plan). The comment atoperators.scala:1040-1044gives the reason:CometExchangeSink.shouldUseShuffleScanonly fires for AQE wrappers (ShuffleQueryStageExec), so a bare non-AQECometShuffleExchangeExecalways serializes as a regular Scan, whateverspark.comet.shuffle.directRead.enabledsays.Profile of one small Comet task, run alone about 800 times, from a 1 ms stack sampler on the task thread plus macOS
sample:org.apache.arrow.vector.ipcmessage objects 6%, JNI releases 5%decodeShuffleBlock)Unsafe.allocateMemory0,Unsafe.freeMemory0,AllocationManager) takes about 22% of the task by itself, because every block carries hundreds of buffers for its 104 columns.sampleattributes 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 inSystem.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
ShuffleQueryStageExec. That removes the JVM leg for every block, and with AQE these same tasks already beat Spark.Additional context
NativeBatchDecoderIterator. This issue is the native-consumer case without AQE, where no JVM code needs the batch at all.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.SELECT /*+ SHUFFLE_HASH(d) */ ... FROM fact f JOIN dim d ON f.key = d.k.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.CometHashJoinTaskTimeBenchmark -- profile=small).mainat 9a6241f with the unmerged perf: read shuffle blocks in 64 KiB pieces instead of throughChannels.newChannel#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