Describe the bug
In Comet's hash-based JVM columnar shuffle, CometBypassMergeSortShuffleWriter, each partition writer flushes whenever it reaches spark.comet.shuffle.jvm.batchSize rows (8192 by default, CometDiskBlockWriter.java#L216-L218). The flush writes that batch to the partition's output file with doSpilling(false), which hands native a throwaway ShuffleWriteMetrics. It then adds the bytes to the task's diskBytesSpilled and never to the shuffle's bytes written (CometDiskBlockWriter.java#L322-L362). Only the final doSpilling(true) in close() counts as bytes written.
For this writer, then, the Spark UI's shuffle write size is roughly one partial batch per partition. "Spill (Disk)" shows almost all of the shuffle output, even with no memory pressure. The map output itself is correct, because MapStatus is built from the file segments.
Steps to reproduce
I found this by reading the code and haven't run it. Any hash-based JVM shuffle with more than one batch per partition should show it. For example, with spark.comet.shuffle.mode=jvm, 10 shuffle partitions and a few million rows, compare the stage's shuffle write size with the map output sizes.
Expected behavior
Every batch written to a partition's output file counts toward shuffle bytes written, which is where Spark's own bypass writer counts it. Nothing is reported as spill unless it went to a separate spill file.
Describe the bug
In Comet's hash-based JVM columnar shuffle,
CometBypassMergeSortShuffleWriter, each partition writer flushes whenever it reachesspark.comet.shuffle.jvm.batchSizerows (8192 by default, CometDiskBlockWriter.java#L216-L218). The flush writes that batch to the partition's output file withdoSpilling(false), which hands native a throwawayShuffleWriteMetrics. It then adds the bytes to the task'sdiskBytesSpilledand never to the shuffle's bytes written (CometDiskBlockWriter.java#L322-L362). Only the finaldoSpilling(true)inclose()counts as bytes written.For this writer, then, the Spark UI's shuffle write size is roughly one partial batch per partition. "Spill (Disk)" shows almost all of the shuffle output, even with no memory pressure. The map output itself is correct, because
MapStatusis built from the file segments.Steps to reproduce
I found this by reading the code and haven't run it. Any hash-based JVM shuffle with more than one batch per partition should show it. For example, with
spark.comet.shuffle.mode=jvm, 10 shuffle partitions and a few million rows, compare the stage's shuffle write size with the map output sizes.Expected behavior
Every batch written to a partition's output file counts toward shuffle bytes written, which is where Spark's own bypass writer counts it. Nothing is reported as spill unless it went to a separate spill file.