Skip to content

JVM columnar shuffle fails with ArrayIndexOutOfBoundsException when spark.shuffle.checksum.enabled=false #6256

Description

@andygrove

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). writeSortedFileNative guards every other access with partitionChecksums.length > 0, but not the store at 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:

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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:shuffleShuffle (JVM and native)bugSomething isn't workingpriority:highCrashes, panics, segfaults, major functional breakage

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions