Skip to content

fix: count every hash-based JVM shuffle batch as shuffle bytes written - #6369

Open
andygrove wants to merge 2 commits into
apache:mainfrom
andygrove:jvm-shuffle-bytes-written-metric
Open

andygrove wants to merge 2 commits into
apache:mainfrom
andygrove:jvm-shuffle-bytes-written-metric

Conversation

@andygrove

@andygrove andygrove commented Sep 29, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6258.

Rationale for this change

In the hash-based JVM columnar shuffle (CometBypassMergeSortShuffleWriter), each partition's CometDiskBlockWriter buffers rows and appends them to the partition's file one Arrow IPC batch at a time. A batch is written in three cases:

  • when it reaches spark.comet.shuffle.jvm.batchSize rows,
  • when a memory allocation fails and the writer spills to free memory,
  • once more in close().

writePartitionedData then concatenates the partition files into the map output.

CometDiskBlockWriter.ArrowIPCWriter.doSpilling(boolean isLast) gave native a throwaway ShuffleWriteMetrics for every batch except the one written in close(), and added those bytes to the task's diskBytesSpilled. As a result:

  • shuffle bytes written only covered the last partial batch of each partition,
  • "Spill (Disk)" showed nearly all of the shuffle output, even with no memory pressure,
  • the write time of the earlier batches was dropped,
  • real spills under memory pressure never reported memoryBytesSpilled.

The map output itself was correct.

A batch that fills up is an ordinary write, as it is in Spark's BypassMergeSortShuffleWriter. A batch written because memory ran out is a spill, but it goes into the same partition file and becomes part of the map output. This writer has no merge phase that would count it as written later, so it has to count as both a spill and a write.

What changes are included in this PR?

CometDiskBlockWriter now reports its metrics as follows:

  • Every batch written to a partition file counts toward shuffle bytes written, records written, and write time. The task's ShuffleWriteMetricsReporter is passed to SpillWriter.doSpilling for every batch.
  • A batch written to relieve memory pressure is also reported as a spill. This covers ArrowIPCWriter.spill, which flushes the requesting writer and, if needed, the task's other partition writers. Its bytes go to diskBytesSpilled, and the memory pages it freed go to memoryBytesSpilled, as in CometShuffleExternalSorter.spill and Spark's ShuffleExternalSorter.spill. So its bytes show up in both the shuffle write size and "Spill (Disk)".
  • Batches written on reaching the batch size (or the internal spark.comet.shuffle.jvm.spillThreshold) and the final batch in close() count only as writes.
  • Write time includes the IO of spilled batches. Spark leaves spill IO out of write time only where spill files are merged later and the merge's IO is counted (ShuffleExternalSorter.writeSortedFile, ExternalSorter). Its bypass writer counts every write to its partition files, and a spilled batch here is written once, into its partition file.
  • The log line for a batch-size write no longer calls it a spill, and it moves from INFO to DEBUG because it fired for every batch of every partition.

The sort-based writer (CometUnsafeShuffleWriter / SpillSorter) writes real spill files and merges them into the output. Its accounting already matches Spark's UnsafeShuffleWriter and is unchanged.

How are these changes tested?

A new test in CometTaskMetricsSuite, JVM hash shuffle counts batch-size writes as shuffle bytes written, not spill, runs a jvm mode shuffle of 20,000 rows into 4 partitions with spark.comet.shuffle.jvm.batchSize=100, so each partition writer writes many batches. It asserts that the plan has a Comet columnar shuffle exchange whose shuffle handle is CometBypassMergeSortShuffleHandle, so the hash-based writer ran. It then checks that:

  • the exchange's shuffleBytesWritten SQL metric and the stages' shuffle write bytes in the status store both equal the total size of the shuffle's .data files,
  • records written equal the input rows,
  • memory and disk spill are both 0.

The memory pressure test in CometDiskBlockWriterSuite, memory pressure spills only writers of the requesting task, now checks the spill and write metrics together:

  • For the task under pressure, diskBytesSpilled equals exactly the bytes in its partition files before they are closed (all written by spills), and memoryBytesSpilled is greater than 0.
  • Bytes written for that task equal the same spilled bytes. After closing, bytes written equal the final file sizes, records written equal the rows inserted, and the final batches leave disk spill unchanged.
  • The other task, which is not under pressure, reports no spill.

I ran the new end-to-end test and the memory pressure test before changing the main code. The end-to-end test failed with 17523 did not equal 316630: the reported bytes written covered about 5.5% of the map output. The memory pressure test's bytes-written check failed because the task reported 0 bytes written after its spills.

With the fix, on the default Spark 4.1 profile, CometDiskBlockWriterSuite (4 tests), CometTaskMetricsSuite (25 tests), and CometShuffleSuite (46 tests) all pass.

…ot spill

CometDiskBlockWriter appends every batch, whether written at the batch
size, under memory pressure, or on close, to the partition file that
CometBypassMergeSortShuffleWriter concatenates into the map output. It
counted only the final batch as shuffle bytes written, reported the
earlier batches as disk spill, and dropped their write time. Count every
batch with the task's shuffle write metrics, as Spark's
BypassMergeSortShuffleWriter does for its partition files, and report
nothing as spill.
@github-actions github-actions Bot added bug Something isn't working area:shuffle Shuffle (JVM and native) labels Sep 29, 2026
…pill

A flush that CometDiskBlockWriter makes to relieve memory pressure now
adds its bytes to the task's diskBytesSpilled and the memory it freed
to memoryBytesSpilled, like a sort-based spill. The flush still counts
as shuffle bytes and records written, since it lands in the partition
file that becomes the map output. Flushes on reaching the batch size
and the final flush on close count only as writes. The log line for a
batch-size flush no longer calls it a spill and moves to debug level.
@andygrove andygrove changed the title fix: count hash-based JVM shuffle batches as shuffle bytes written, not spill fix: count every hash-based JVM shuffle batch as shuffle bytes written Sep 29, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:shuffle Shuffle (JVM and native) bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Hash-based JVM columnar shuffle reports its output as disk spill instead of bytes written

1 participant