Skip to content

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

Description

@andygrove

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:shuffleShuffle (JVM and native)bugSomething isn't workingpriority:mediumFunctional bugs, performance regressions, broken features

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions