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]

Reply via email to