andygrove opened a new issue, #6256:
URL: https://github.com/apache/datafusion-comet/issues/6256
### Describe the bug
With `spark.shuffle.checksum.enabled=false`, every task of a Comet JVM
columnar shuffle that uses the sort-based writer fails:
```
java.lang.ArrayIndexOutOfBoundsException: Index 0 out of bounds for length 0
at
org.apache.spark.shuffle.sort.SpillSorter.writeSortedFileNative(SpillSorter.java:290)
at
org.apache.spark.shuffle.sort.CometShuffleExternalSorter.closeAndGetSpills(CometShuffleExternalSorter.java:337)
at
org.apache.spark.sql.comet.execution.shuffle.CometUnsafeShuffleWriter.closeAndWriteOutput(CometUnsafeShuffleWriter.java:309)
at
org.apache.spark.sql.comet.execution.shuffle.CometUnsafeShuffleWriter.write(CometUnsafeShuffleWriter.java:255)
```
With checksums disabled, `createPartitionChecksums` returns an empty array
([CometShuffleChecksumSupport.java#L28-L41](https://github.com/apache/datafusion-comet/blob/bc4be39964cbe9cdb5f2a949740a8164e6b5755b/spark/src/main/java/org/apache/spark/shuffle/comet/CometShuffleChecksumSupport.java#L28-L41)).
`writeSortedFileNative` guards every other access with
`partitionChecksums.length > 0`, but not the store at
[SpillSorter.java#L290](https://github.com/apache/datafusion-comet/blob/bc4be39964cbe9cdb5f2a949740a8164e6b5755b/spark/src/main/java/org/apache/spark/shuffle/sort/SpillSorter.java#L290).
That store runs every time the sorted records move to the next partition.
### Steps to reproduce
With `spark.shuffle.checksum.enabled=false` in the SparkConf:
```scala
spark.conf.set("spark.comet.shuffle.mode", "jvm")
spark.range(0, 100000, 1, 4)
.selectExpr("id", "cast(id as string) as s")
.repartition(300, col("id"))
.count()
```
300 partitions is above `spark.shuffle.sort.bypassMergeThreshold`, so Comet
uses the sort-based writer.
### Expected behavior
The shuffle succeeds with checksums disabled.
### Additional context
The unguarded store has been there since at least #3081, which moved
`SpillSorter` into its own class, so 1.0.0 and 1.1.0 are affected.
`SpillSorterSuite` always passes a non-empty checksum array, which is why no
test caught this.
--
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]