JVM Shuffle#

This document describes Comet’s JVM-based columnar shuffle implementation (CometColumnarShuffle), which writes shuffle data in Arrow IPC format using JVM code with native encoding. For the fully native alternative, see Native Shuffle.

Overview#

Comet provides two shuffle implementations:

  • CometNativeShuffle (CometExchange): Fully native shuffle using Rust. Takes columnar input directly from Comet native operators and performs partitioning in native code.

  • CometColumnarShuffle (CometColumnarExchange): JVM-based shuffle that operates on rows internally, buffers UnsafeRows in memory pages, and uses native code (via JNI) to encode them to Arrow IPC format. Uses Spark’s partitioner for partition assignment. Can accept either row-based or columnar input (columnar input is converted to rows via ColumnarToRowExec).

The JVM shuffle is selected via CometShuffleDependency.shuffleType.

When JVM Shuffle is Used#

JVM shuffle (CometColumnarExchange) is used instead of native shuffle (CometExchange) in the following cases:

  1. Shuffle mode is explicitly set to “jvm”: When spark.comet.shuffle.mode is set to jvm.

  2. Child plan is not a Comet native operator: When the child plan is a Spark row-based operator (not a CometPlan), JVM shuffle is used, since native shuffle requires columnar input from Comet operators. The exception is spark.comet.convert.shuffleInput.enabled, which converts the child’s rows to Arrow with CometSparkToColumnarExec so that native shuffle can take the shuffle instead. See When Native Shuffle is Used.

  3. Unsupported partition key types: RangePartitioning keys must be primitive, so a complex range key always falls back here. HashPartitioning keys must be primitive only by default: setting spark.comet.shuffle.native.partitioning.hash.nested.enabled to true keeps struct and array keys, and map keys on Spark 4.0 and later, on the native path. The config defaults to false, so a complex hash key falls back to JVM columnar shuffle unless it is enabled. See Supported partition key types for the exact rules. Complex types are fully supported as data columns in both implementations.

Input Handling#

Spark Row-Based Input#

When the child plan is a Spark row-based operator, CometColumnarExchange calls child.execute() which returns an RDD[InternalRow]. The rows flow directly to the JVM shuffle writers.

Comet Columnar Input#

When the child plan is a Comet native operator (e.g., CometHashAggregate) but JVM shuffle is selected (due to shuffle mode setting or unsupported partitioning), CometColumnarExchange still calls child.execute(). Comet operators implement doExecute() by wrapping themselves with ColumnarToRowExec:

// In CometExec base class
override def doExecute(): RDD[InternalRow] =
  ColumnarToRowExec(this).doExecute()

This means the data path becomes:

Comet Native (columnar) → ColumnarToRowExec → rows → JVM Shuffle → Arrow IPC → columnar

This is less efficient than native shuffle which avoids the columnar-to-row conversion:

Comet Native (columnar) → Native Shuffle → Arrow IPC → columnar

Why Use Spark’s Partitioner?#

JVM shuffle uses row-based input so it can leverage Spark’s existing partitioner infrastructure (partitioner.getPartition(key)). This allows Comet to support all of Spark’s partitioning schemes without reimplementing them in Rust. Native shuffle, by contrast, serializes the partitioning scheme to protobuf and implements the partitioning logic in native code.

Architecture#

┌─────────────────────────────────────────────────────────────────────────┐
│                         CometShuffleManager                             │
│  - Extends Spark's ShuffleManager                                       │
│  - Routes to appropriate writer/reader based on ShuffleHandle type      │
└─────────────────────────────────────────────────────────────────────────┘
                                    │
           ┌────────────────────────┼────────────────────────┐
           ▼                        ▼                        ▼
┌─────────────────────┐  ┌─────────────────────┐  ┌─────────────────────┐
│ CometBypassMerge-   │  │ CometUnsafe-        │  │ CometNative-        │
│ SortShuffleWriter   │  │ ShuffleWriter       │  │ ShuffleWriter       │
│ (hash-based)        │  │ (sort-based)        │  │ (fully native)      │
└─────────────────────┘  └─────────────────────┘  └─────────────────────┘
           │                        │
           ▼                        ▼
┌─────────────────────┐  ┌─────────────────────┐
│ CometDiskBlock-     │  │ CometShuffleExternal│
│ Writer              │  │ Sorter              │
└─────────────────────┘  └─────────────────────┘
           │                        │
           └────────────┬───────────┘
                        ▼
              ┌─────────────────────┐
              │ SpillWriter         │
              │ (native encoding    │
              │  via JNI)           │
              └─────────────────────┘

Key Classes#

Shuffle Manager#

Class

Location

Description

CometShuffleManager

.../shuffle/CometShuffleManager.scala

Entry point. Extends Spark’s ShuffleManager. Selects writer/reader based on handle type. Delegates non-Comet shuffles to SortShuffleManager.

CometShuffleDependency

.../shuffle/CometShuffleDependency.scala

Extends ShuffleDependency. Contains shuffleType (CometColumnarShuffle or CometNativeShuffle) and schema info.

Shuffle Handles#

Handle

Writer Strategy

CometBypassMergeSortShuffleHandle

Hash-based: one file per partition, merged at end

CometSerializedShuffleHandle

Sort-based: records sorted by partition ID, single output

CometNativeShuffleHandle

Fully native shuffle

Selection logic in CometShuffleManager.shouldBypassMergeSort():

  • Uses bypass if partitions < threshold AND partitions × cores ≤ max threads

  • Otherwise uses sort-based to avoid OOM from many concurrent writers

Writers#

Class

Location

Description

CometBypassMergeSortShuffleWriter

.../shuffle/CometBypassMergeSortShuffleWriter.java

Hash-based writer. Creates one CometDiskBlockWriter per partition. Supports async writes.

CometUnsafeShuffleWriter

.../shuffle/CometUnsafeShuffleWriter.java

Sort-based writer. Uses CometShuffleExternalSorter to buffer and sort records, then merges spill files.

CometDiskBlockWriter

.../shuffle/CometDiskBlockWriter.java

Buffers rows in memory pages for a single partition. Spills to disk via native encoding. Used by bypass writer.

CometShuffleExternalSorter

.../shuffle/sort/CometShuffleExternalSorter.java

Buffers records across all partitions, sorts by partition ID, spills sorted data. Used by unsafe writer.

SpillWriter

.../shuffle/SpillWriter.java

Base class for spill logic. Manages memory pages and calls Native.writeSortedFileNative() for Arrow IPC encoding.

Reader#

Class

Location

Description

CometBlockStoreShuffleReader

.../shuffle/CometBlockStoreShuffleReader.scala

Fetches shuffle blocks via ShuffleBlockFetcherIterator. Decodes Arrow IPC to ColumnarBatch.

NativeBatchDecoderIterator

.../shuffle/NativeBatchDecoderIterator.scala

Reads compressed Arrow IPC batches from input stream. Calls Native.decodeShuffleBlock() via JNI.

Data Flow#

Write Path#

  1. ShuffleWriteProcessor calls CometShuffleManager.getWriter()

  2. Writer receives Iterator[Product2[K, V]] where V is UnsafeRow

  3. Rows are serialized and buffered in off-heap memory pages

  4. When memory threshold or batch size is reached, SpillWriter.doSpilling() is called

  5. Native code (Native.writeSortedFileNative()) converts rows to Arrow arrays and writes IPC format

  6. For bypass writer: partition files are concatenated into final output

  7. For sort writer: spill files are merged

Read Path#

  1. CometBlockStoreShuffleReader.read() creates ShuffleBlockFetcherIterator

  2. For each block, NativeBatchDecoderIterator reads the IPC stream

  3. Native code (Native.decodeShuffleBlock()) decompresses and decodes to Arrow arrays

  4. Arrow FFI imports arrays as ColumnarBatch

This is the path taken when the consumer is a Spark operator. When a Comet native operator consumes the shuffle output and spark.comet.shuffle.directRead.enabled is set, steps 2 through 4 are skipped: the raw compressed blocks are handed to native code, which decodes them inside the plan. JVM shuffle writes the same Arrow IPC block format as native shuffle, so direct read applies to both. See Direct Read for the mechanism.

Memory Management#

  • CometShuffleMemoryAllocator.getInstance returns the allocator for the task. Pages are always Unsafe-allocated, because their addresses are handed to native code, but what bounds them depends on Spark’s memory mode.

  • Off-heap mode gets CometUnifiedShuffleMemoryAllocator, an ordinary Spark MemoryConsumer drawing from spark.memory.offHeap.size. When it cannot acquire a page the writer spills to disk. An allocation that loses the task’s entry in Spark’s execution pool while it waits (SPARK-59444) is retried a few times before it counts as a page that could not be acquired.

  • On-heap mode gets CometUnboundedShuffleMemoryAllocator, which keeps no budget and so never refuses a page. Nothing bounds these allocations, and memory pressure never triggers a spill. That mode exists only so the Spark SQL tests can run against Comet. See Memory Management.

  • Row count still triggers spilling in either mode. CometDiskBlockWriter spills at min(spark.comet.shuffle.jvm.spillThreshold, batch size), where the batch size is spark.comet.shuffle.jvm.batchSize capped at spark.comet.batchSize. CometShuffleExternalSorter spills at spark.comet.shuffle.jvm.spillThreshold alone, which defaults to Int.MaxValue, so on the sort path that trigger is effectively off by default.

  • CometDiskBlockWriter coordinates spilling across all partition writers (largest first)

Configuration#

Config

Description

spark.comet.shuffle.jvm.batchSize

Rows per Arrow batch

spark.comet.shuffle.jvm.spillThreshold

Row count threshold for spill

spark.comet.shuffle.compression.codec

Compression codec, lz4 by default

spark.comet.shuffle.directRead.enabled

Decode blocks natively on read, true by default