Conversation
…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.
…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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Which issue does this PR close?
Closes #6258.
Rationale for this change
In the hash-based JVM columnar shuffle (
CometBypassMergeSortShuffleWriter), each partition'sCometDiskBlockWriterbuffers rows and appends them to the partition's file one Arrow IPC batch at a time. A batch is written in three cases:spark.comet.shuffle.jvm.batchSizerows,close().writePartitionedDatathen concatenates the partition files into the map output.CometDiskBlockWriter.ArrowIPCWriter.doSpilling(boolean isLast)gave native a throwawayShuffleWriteMetricsfor every batch except the one written inclose(), and added those bytes to the task'sdiskBytesSpilled. As a result: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?
CometDiskBlockWriternow reports its metrics as follows:ShuffleWriteMetricsReporteris passed toSpillWriter.doSpillingfor every batch.ArrowIPCWriter.spill, which flushes the requesting writer and, if needed, the task's other partition writers. Its bytes go todiskBytesSpilled, and the memory pages it freed go tomemoryBytesSpilled, as inCometShuffleExternalSorter.spilland Spark'sShuffleExternalSorter.spill. So its bytes show up in both the shuffle write size and "Spill (Disk)".spark.comet.shuffle.jvm.spillThreshold) and the final batch inclose()count only as writes.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 sort-based writer (
CometUnsafeShuffleWriter/SpillSorter) writes real spill files and merges them into the output. Its accounting already matches Spark'sUnsafeShuffleWriterand 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 ajvmmode shuffle of 20,000 rows into 4 partitions withspark.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 isCometBypassMergeSortShuffleHandle, so the hash-based writer ran. It then checks that:shuffleBytesWrittenSQL metric and the stages' shuffle write bytes in the status store both equal the total size of the shuffle's.datafiles,The memory pressure test in
CometDiskBlockWriterSuite,memory pressure spills only writers of the requesting task, now checks the spill and write metrics together:diskBytesSpilledequals exactly the bytes in its partition files before they are closed (all written by spills), andmemoryBytesSpilledis greater than 0.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), andCometShuffleSuite(46 tests) all pass.