andygrove opened a new pull request, #6369: URL: https://github.com/apache/datafusion-comet/pull/6369
## 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 writes them to the partition's file one Arrow IPC batch at a time. It writes a batch when it reaches `spark.comet.shuffle.jvm.batchSize` rows (or the internal spill threshold), when memory pressure forces a flush, and once more in `close()`. Every batch is appended to the same partition file, and `writePartitionedData` concatenates the partition files into the map output. `CometDiskBlockWriter.ArrowIPCWriter.doSpilling(boolean isLast)` followed Spark's `ShuffleExternalSorter.writeSortedFile`, where non-final data goes to separate spill files that are merged later. For every batch except the one written in `close()`, it passed native a throwaway `ShuffleWriteMetrics`, copied the record count to the real metrics, and added the 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 from shuffle write time. The map output itself was correct. This writer has no separate spill files. Under memory pressure it writes its buffered rows to the same partition file, so those batches are part of the map output as well. Every batch should therefore be counted the way Spark's `BypassMergeSortShuffleWriter` counts writes to its per-partition files: as shuffle bytes written, records written, and write time, with nothing reported as spill. The sort-based writer (`CometUnsafeShuffleWriter` / `SpillSorter`) does write real spill files and merges them into the output. Its accounting (spill files as disk spill, merged output as bytes written) already matches Spark's `UnsafeShuffleWriter`, so this PR leaves it unchanged. ## What changes are included in this PR? - `CometDiskBlockWriter.ArrowIPCWriter.doSpilling` no longer takes an `isLast` flag. It writes every batch with the task's `ShuffleWriteMetricsReporter`, so `SpillWriter.doSpilling` adds the batch's bytes, records, and write time to the shuffle write metrics, and nothing goes to `diskBytesSpilled`. Records written were already counted for every batch, and each batch is still counted once. - `CometDiskBlockWriterSuite`: the memory pressure test previously asserted that the early flushes were reported as disk spill. It now checks that they are reported as bytes written that match the partition file sizes, that records written equal the rows inserted, and that no disk spill is reported. - `CometTaskMetricsSuite`: a new end-to-end test, described below. ## How are these changes tested? The new test `JVM hash shuffle reports every written batch as shuffle bytes, not spill` in `CometTaskMetricsSuite` 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 on disk, - records written equal the input rows, - memory and disk spill are both 0. I ran both tests before changing the main code, and both failed. The new test failed with `17523 did not equal 316630`: the reported bytes written covered about 5.5% of the map output. The updated `CometDiskBlockWriterSuite` test failed because task A reported 0 bytes written after its memory-pressure flushes. With the fix, on the default Spark 4.1 profile, `CometDiskBlockWriterSuite` (4 tests), `CometTaskMetricsSuite` (25 tests), and `CometShuffleSuite` (46 tests) all pass. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
