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]
