Repository navigation
Conversation
…=false With shuffle checksums disabled, the partition checksum array is empty. SpillSorter.writeSortedFileNative stored a checksum into it whenever the sorted records moved to the next partition, so every task of the sort-based JVM columnar shuffle writer failed with an ArrayIndexOutOfBoundsException. Guard that store like the other accesses. SpillWriter.doSpilling also replaced its -1 "no checksum" sentinel with the Long.MIN_VALUE the native writer returns when it computes no checksum, so every later write by the same writer computed an unused Adler32 checksum. Only keep the returned checksum when checksums are enabled.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Disabling shuffle checksums caused partition-switch failures and unintentionally enabled checksum computation after the first write.
- Design approach: Guard the checksum-array store and preserve the writer’s disabled-checksum sentinel.
- Correctness / compatibility analysis: The production changes match Spark’s empty-checksum-array contract across supported versions 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. One introduced build blocker remains: the new test resolves
SpillInfoincorrectly under Scala 2.12. - Key design decisions: Reuses existing state and conventions without additional abstractions, allocations or JNI calls. Avoiding unintended checksum computation removes demonstrated unnecessary work.
- Implementation sketch: Two production guards, three regression tests covering partition transitions and repeated spills, and suite registration in both platform workflows.
- Behavioral changes worth calling out: JVM shuffle can retain disabled checksums across partitions and spills. Enabled checksum updates retain their existing behavior.
- Suggested improvements: Alias Comet’s
SpillInfoin the new test and use that alias for both allocations, then rerun the affected builds.
Reviewed full SHA 207fab1faa05f7ab0cce7e07254a0e4a5a1e43d4. Covered all seven files in the PR diff against base 22a07067bda272ed68adc67733cbaa7462a13c85, using merge base 65b334bfdd25196091a42d773bc3e61d19212e52. Confirmed non-draft status and read the discussion snapshot and live comments. No existing review concerns were present.
Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-memory-pr, and review-comet-ffi-pr.
Exact-head CI: Linux Spark 4.1 suites passed, including all 515 shuffle tests and the three new regressions. Native build, Rust tests, TPC-H/TPC-DS verification and Spark 4.0 build/lint checks passed. Spark 3.4/3.5 build checks and both Celeborn compatibility jobs failed at the new test’s type mismatch, leaving Required Checks red.
Validation limits: A bounded local compiler reproduction confirmed the failure and alias fix. No full local Comet build or runtime suite was run. macOS, Spark SQL and Iceberg suites were skipped in exact-head CI.
| test("write sorted file across partitions with shuffle checksums disabled") { | ||
| // With spark.shuffle.checksum.enabled=false the sorter gets an empty checksum array and no | ||
| // checksum algorithm. | ||
| val spills = new java.util.LinkedList[SpillInfo]() |
There was a problem hiding this comment.
[P1] Alias Comet’s SpillInfo to restore Scala 2.12 builds. When compiling the supported Spark 3.4 or 3.5 profile, Spark’s org.apache.spark.shuffle.sort.SpillInfo in this suite’s package hides the newly imported Comet class. This allocation therefore produces the wrong LinkedList type, and passing it to createSpillSorter fails compilation instead of running the regression test. The constructor at line 294 has the same collision. This blocks both profiles’ test compilation and both Celeborn compatibility jobs. Use an import alias such as SpillInfo => CometSpillInfo at both allocations, or fully qualify the Comet type.
Evidence: Exact-head CI reports found: java.util.LinkedList[org.apache.spark.shuffle.sort.SpillInfo], required: java.util.LinkedList[org.apache.spark.sql.comet.execution.shuffle.SpillInfo] at line 276, plus the hidden-import warning. See Spark 3.5 job https://github.com/apache/datafusion-comet/actions/runs/36499297541/job/109186916881 and Spark 3.4 job https://github.com/apache/datafusion-comet/actions/runs/36499297541/job/109186916902. A local JDK 17 compiler probe using Spark 3.5.9 and the checkout’s actual Comet SpillInfo reproduced both mismatches with Scala 2.12.18. Aliasing both uses compiled successfully. Original and aliased forms both compiled with Scala 2.13.17.
comphead
left a comment
There was a problem hiding this comment.
Thanks for the careful fix and for auditing the other checksum uses. Guarding on partitionChecksums.length > 0 matches Spark's ShuffleExternalSorter, and Spark's IndexShuffleBlockResolver treats an empty array as checksums disabled (checksums.nonEmpty), in 3.4.3 through 4.1.2. I didn't spot another unguarded use. The comments below are small test and cleanup suggestions.
| (0 until 5).foreach(_ => writer.insertRow(toUnsafe(InternalRow(new Array[Byte](16))), 0)) | ||
| assert(writer.close().length > 0) | ||
| assert(writer.getOutputRecords == 5) | ||
| assert(writer.getChecksum == -1) |
There was a problem hiding this comment.
Would it make sense to also run this scenario with checksums enabled? That is the default, and this PR changes the assignment in doSpilling that the enabled path relies on, but I couldn't find a test that checks a checksum value. With writer.setChecksum(Long.MIN_VALUE) and writer.setChecksumAlgo("adler32") before the inserts, I'd expect getChecksum to equal a one-shot java.util.zip.Adler32 over the partition file (I haven't run this). Parameterizing this test over both settings would keep the test count flat.
| .range(0, 100000, 1, 4) | ||
| .selectExpr("id", "cast(id as string) as s") | ||
| .repartition(numPartitions, col("id")) | ||
| checkCometExchange(df, 1, false) |
There was a problem hiding this comment.
The comment says 300 partitions use the sort-based writer. Would it be worth asserting that, so the test can't silently stop covering it if a threshold changes? checkCometExchange returns the exchanges, so something like assert(exchanges.head.shuffleDependency.shuffleHandle.getClass.getName.contains("CometSerializedShuffleHandle") == (numPartitions > 200)) should work (I haven't run it). The similar checks in Columnar shuffle for large shuffle partition number discard the boolean, so as written they don't assert anything.
|
|
||
| // Store the checksum for the current partition. | ||
| partitionChecksums[currentPartition] = getChecksum(); | ||
| if (partitionChecksums.length > 0) { |
There was a problem hiding this comment.
Nit: the partition-switch block here and the final-partition block below are near copies, and the final one already had this guard, so they drifted apart. Would a small private helper that sets the checksum, calls doSpilling, records the length and stores the checksum keep the guard in one place? Happy to leave that for a follow-up if you'd rather keep this fix small for backporting.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Disabling shuffle checksums caused partition-switch failures and unintentionally enabled checksum computation after the first write.
- Design approach: Guard the checksum-array store and preserve the writer’s disabled-checksum sentinel.
- Correctness / compatibility analysis: The production guards match Spark’s empty-array contract across supported versions 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The previously reported Scala 2.12 compilation blocker remains unresolved.
- Key design decisions: Reuses existing state without additional abstractions, allocations or JNI calls. Prevents unnecessary checksum computation without changing enabled-checksum updates.
- Implementation sketch: Two production guards, three regression tests covering partition transitions and repeated spills, and registration in both platform workflows.
- Behavioral changes worth calling out: JVM shuffle retains disabled checksums across partitions and spills. Partitioning, output format and memory ownership remain unchanged.
- Suggested improvements: Address the existing
SpillInfoimport collision using an alias at both allocations, then rerun the affected builds. No additional introduced P1/P2 issues found within this review.
Reviewed full SHA 207fab1faa05f7ab0cce7e07254a0e4a5a1e43d4, covering all seven files in the full PR diff against base 22a07067bda272ed68adc67733cbaa7462a13c85, using merge base 65b334bfdd25196091a42d773bc3e61d19212e52. Confirmed non-draft status and read existing reviews, comments and threads.
Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-memory-pr and review-comet-ffi-pr.
Existing blocker: The reported P1 remains reproducible. Scala 2.12 resolves the new test’s SpillInfo references to Spark’s class instead of Comet’s. A bounded local compiler reproduction with JDK 17, Scala 2.12.18, Spark 3.5.9 and the checkout’s Comet class reproduced both mismatches. Aliasing both uses compiled successfully. This is not duplicated as a new finding.
Exact-head CI: All 38 checks reference the requested head. The build jobs tested GitHub’s merge with the requested base. Linux Spark 4.1 passed all 515 shuffle tests, including the three new regressions. Other Linux 4.1 suites, native/Rust checks, TPC-H/TPC-DS verification and Spark 4.0 build checks passed. Spark 3.4/3.5 and both Celeborn jobs failed on the existing type collision, leaving Required Checks red.
Validation limits: No full local Comet build or runtime suite was run. Runtime CI coverage was limited to the default Spark profile. macOS, Spark SQL and Iceberg suites were skipped.
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 exceedsspark.shuffle.sort.bypassMergeThreshold) fails withArrayIndexOutOfBoundsException: Index 0 out of bounds for length 0.CometShuffleChecksumSupport.createPartitionChecksumsreturns an empty array when checksums are disabled.SpillSorter.writeSortedFileNativecheckspartitionChecksums.length > 0before 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.
SpillWriteruseschecksum == -1to mean "no checksum", butdoSpillingalways overwriteschecksumwith the value returned by the nativewriteSortedFileNative, which isLong.MIN_VALUEwhen no checksum was computed. From its second write on, a writer therefore passeschecksumEnabled = true, and the native code computes an Adler32 checksum that nothing reads.SpillSorteris reused across partitions and spills, andCometDiskBlockWriteracross 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.
CometBypassMergeSortShuffleWriterguardssetChecksumand collects checksums with a loop bounded bypartitionChecksums.length.CometShuffleExternalSorterandCometUnsafeShuffleWriteronly pass the (possibly empty) array to Spark'scommitAllPartitionsandtransferMapSpillFile, which treat an empty array as checksums disabled, just as for Spark's ownUnsafeShuffleWriter.What changes are included in this PR?
SpillSorter.writeSortedFileNative: guard the per-partition checksum store withpartitionChecksums.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.CometShuffleChecksumDisabledSuiteis registered in theshufflegroup 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 beforeclose(), then checks thatgetChecksum()still reports no checksum.CometShuffleChecksumDisabledSuiteinCometColumnarShuffleSuite.scala.spark.shuffle.checksum.enabledis a core setting, so the suite sets it insparkConf, asCometShuffleEncryptionSuitedoes for encryption. It runs the query from the issue withspark.comet.shuffle.mode=jvmfor 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
SpillSorterSuitetest and the 300-partition case of the end-to-end test failed with theArrayIndexOutOfBoundsExceptionatSpillSorter.java:290from the issue. TheCometDiskBlockWriterSuitetest failed with2642559282 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) andCometShuffleSuite(46, JVM columnar shuffle with checksums enabled and forced spills) all pass.