fix: validate spark.comet.shuffle.jvm.batchSize and spark.comet.exec.memoryPool - #6287
Conversation
…memoryPool spark.comet.shuffle.jvm.batchSize was only checked to be no larger than spark.comet.batchSize, so 0 was accepted. The JVM shuffle writer then passed 0 to process_sorted_row_partition, whose loop never advanced, and because the loop runs inside a JNI call the task hung and could not be killed. The value must now be positive. spark.comet.exec.memoryPool had no check, so a misspelled value, or one that differed only in case, reached native code, and in off-heap mode every task failed with "Unsupported memory pool type" when it created a native plan. The value is now lowercased and must be fair_unified or greedy_unified, the same way spark.comet.shuffle.mode is handled. Closes apache#6259
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Zero shuffle batch sizes could stall the native writer. Mixed-case or unknown memory-pool names reached native code and failed there.
- Design approach: Reuse the existing config builders to require positive batch sizes and normalize pool names with
Locale.ROOTbefore validating them. - Correctness / compatibility analysis: Traced both JVM shuffle writers and
getMemoryConfigthrough their native consumers. Checked Spark configuration sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. The accepted pool names match native support. No introduced P1/P2 issues found within this review. - Key design decisions: Defaults, deprecated aliases, the existing batch-size upper bound, and on-heap behavior remain unchanged. The implementation adds no new abstraction. Additional work occurs at existing config reads, with no added per-row processing.
- Implementation sketch: Three production builder additions and two regression tests in
CometConfSuitecover rejection and normalization. - Behavioral changes worth calling out: Invalid values fail when read, usually in executor tasks, rather than when set. Mixed-case supported pool names now work. The separately tracked initialization problem in #6286 predates this PR and is unchanged.
- Suggested improvements: None meeting the P1/P2 reporting threshold.
Reviewed full SHA 54676eb52afa793f4fb23bcc3463aae1fbf05907 against requested base 634e37d08b02da4f9fcbe9dea3606dfd884644c7. Verified the complete PR scope using Git ancestry: only the two configuration files changed since the fork. Seven additional tip-to-tip differences originate from later commits on the base branch, and those files are unchanged on this branch. No prerequisite commits were omitted.
Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-shuffle-pr. Also consulted review-comet-iceberg-write-pr during ancestry verification. The PR remains non-draft. Snapshot and live discussion checks found no reviews, comments, or threads.
Exact-head CI: 24 successful checks, including Required Checks, Linux Spark 4.1 suites, Rust tests, and TPC-H/TPC-DS verification. Fourteen checks were skipped, including macOS, Spark SQL, and Iceberg integration suites. No failed or cancelled checks.
Validation: Isolated compilation against Spark 4.1.3 passed all 18 CometConfSuite tests and 30 additional configuration checks covering defaults, boundaries, aliases, thread-local reads, and Turkish-locale casing. Both new regression tests fail against the supplied base. This used current configuration sources with a test-only Utils.stringToSeq stand-in. No full Maven/native build or local end-to-end shuffle execution was performed. Project files remain unchanged.
Which issue does this PR close?
Closes #6259.
Rationale for this change
spark.comet.shuffle.jvm.batchSize=0passed the only check the entry had, that it is no larger thanspark.comet.batchSize. The JVM shuffle writer then hands 0 toprocess_sorted_row_partition, whose loop never advances, and since the loop runs inside a JNI call the task hangs and can't be killed.spark.comet.exec.memoryPoolhad no check at all, so a misspelled value, or one that differs only in case, reached native code, and in off-heap mode every task failed withUnsupported memory pool typewhen it created a native plan.What changes are included in this PR?
spark.comet.shuffle.jvm.batchSizemust be positive. The new check runs ahead of the existing one.spark.comet.exec.memoryPoolis lowercased and must befair_unifiedorgreedy_unified, the same wayspark.comet.shuffle.modeis handled.getMemoryConfigreads it through the entry, so native code gets the lowercased value.Comet checks a config when it is read, and both of these are read in the executor task, so a bad value still fails the task. It now fails straight away, with an error that names the config and the values it accepts. Rejecting a value when it is set would mean checking Comet's configs on the driver, which is a bigger change than this.
The existing check on
spark.comet.shuffle.jvm.batchSizehas a problem of its own, filed as #6286. It readsspark.comet.batchSizewhileCometConfis being initialized, so a batch size below 8192 stopsCometConffrom initializing on executors. This PR leaves that check as it is.How are these changes tested?
CometConfSuite. The first checks that 0 and -1 are rejected as the batch size and 1 is accepted. The second checks that the pool type is lowercased and that a misspelling is rejected with the valid values in the message. Both fail on main.Native.writeSortedFileNativeon both the bypass-merge path (4 partitions) and the sort path (300 partitions), andGREEDY_UNIFIED,Fair_Unifiedandfaireach failed every task withUnsupported memory pool type. With this change the batch size of 0 fails the task within a second, the two case variants run natively with answers that match Spark, andfairfails withshould be one of fair_unified, greedy_unified, but was fair. A test of the hang itself would hang CI rather than fail, so that one isn't committed.