andygrove opened a new pull request, #6361:
URL: https://github.com/apache/datafusion-comet/pull/6361

   ## Which issue does this PR close?
   
   Closes #6256.
   
   ## Rationale for this change
   
   With `spark.shuffle.checksum.enabled=false`, every map task of a JVM 
columnar shuffle that uses the sort-based writer (`CometUnsafeShuffleWriter`, 
chosen when the partition count exceeds 
`spark.shuffle.sort.bypassMergeThreshold`) fails with 
`ArrayIndexOutOfBoundsException: Index 0 out of bounds for length 0`.
   
   `CometShuffleChecksumSupport.createPartitionChecksums` returns an empty 
array when checksums are disabled. `SpillSorter.writeSortedFileNative` checks 
`partitionChecksums.length > 0` before every access except the store that runs 
when the sorted records move on to the next partition, so the first partition 
switch indexes into the empty array.
   
   Auditing the other JVM shuffle paths turned up a second problem in the same 
checksum-disabled path. `SpillWriter` uses `checksum == -1` to mean "no 
checksum", but `doSpilling` always overwrites `checksum` with the value 
returned by the native `writeSortedFileNative`, which is `Long.MIN_VALUE` when 
no checksum was computed. From its second write on, a writer therefore passes 
`checksumEnabled = true`, and the native code computes an Adler32 checksum that 
nothing reads. `SpillSorter` is reused across partitions and spills, and 
`CometDiskBlockWriter` across the spills of a partition, so with checksums 
disabled the JVM shuffle still checksummed almost all of its output. Once the 
store above is guarded, this is the only remaining effect, so it is fixed here 
too.
   
   The other uses of the checksum array are already safe. 
`CometBypassMergeSortShuffleWriter` guards `setChecksum` and collects checksums 
with a loop bounded by `partitionChecksums.length`. 
`CometShuffleExternalSorter` and `CometUnsafeShuffleWriter` only pass the 
(possibly empty) array to Spark's `commitAllPartitions` and 
`transferMapSpillFile`, which treat an empty array as checksums disabled, just 
as for Spark's own `UnsafeShuffleWriter`.
   
   ## What changes are included in this PR?
   
   - `SpillSorter.writeSortedFileNative`: guard the per-partition checksum 
store with `partitionChecksums.length > 0`, like the other accesses.
   - `SpillWriter.doSpilling`: only keep the checksum returned by the native 
code when checksums are enabled, so a writer that starts with checksums 
disabled keeps them disabled for all of its writes.
   - New tests (below). The new `CometShuffleChecksumDisabledSuite` is 
registered in the `shuffle` group of both PR workflows.
   
   ## How are these changes tested?
   
   - `SpillSorterSuite`: the new test "write sorted file across partitions with 
shuffle checksums disabled" writes a sorted file across three partitions with 
an empty checksum array and no checksum algorithm, which is what the sorter 
receives when checksums are disabled. The existing tests always passed a 
non-empty array and never wrote a file, which is why nothing caught this.
   - `CometDiskBlockWriterSuite`: the new test "a writer computes no checksum 
when shuffle checksums are disabled" makes a writer with no checksum set spill 
twice before `close()`, then checks that `getChecksum()` still reports no 
checksum.
   - New `CometShuffleChecksumDisabledSuite` in 
`CometColumnarShuffleSuite.scala`. `spark.shuffle.checksum.enabled` is a core 
setting, so the suite sets it in `sparkConf`, as `CometShuffleEncryptionSuite` 
does for encryption. It runs the query from the issue with 
`spark.comet.shuffle.mode=jvm` for 10 partitions (bypass merge sort writer) and 
300 partitions (sort-based writer). Each runs with the default spill threshold 
and with a threshold of 2000 rows that makes every map task spill, so the 
sort-based writer also merges spill files. The test asserts that the plan has a 
Comet columnar shuffle exchange and that the result matches Spark.
   
   I wrote the tests before the fix and ran them on the default Spark 4.1 
profile. Without the fix, all three failed. The `SpillSorterSuite` test and the 
300-partition case of the end-to-end test failed with the 
`ArrayIndexOutOfBoundsException` at `SpillSorter.java:290` from the issue. The 
`CometDiskBlockWriterSuite` test failed with `2642559282 did not equal -1`, 
which shows the writer computed a checksum. With the fix, on the same profile, 
`SpillSorterSuite` (11 tests), `CometDiskBlockWriterSuite` (5), 
`CometShuffleChecksumDisabledSuite` (1) and `CometShuffleSuite` (46, JVM 
columnar shuffle with checksums enabled and forced spills) 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