Skip to content

fix: cap spark.comet.shuffle.jvm.batchSize at spark.comet.batchSize where it is read - #6366

Queued
andygrove wants to merge 2 commits into
apache:mainfrom
andygrove:jvm-shuffle-batch-size-check-at-read
Queued

andygrove wants to merge 2 commits into
apache:mainfrom
andygrove:jvm-shuffle-batch-size-check-at-read

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6286.

Rationale for this change

spark.comet.shuffle.jvm.batchSize had a validator that compared the value with COMET_BATCH_SIZE.get(). TypedConfigBuilder.createWithDefault runs an entry's validators on its default, and it does that inside CometConf's static initializer, so the check read SQLConf.get at whatever moment CometConf was first loaded. Nothing loads CometConf when an executor starts, so the first load happens inside a task, where SQLConf.get holds the session's confs. With spark.comet.batchSize below 8192, the default of 8192 failed the check, CometConf threw ExceptionInInitializerError, and every later task on that executor failed with NoClassDefFoundError: Could not initialize class org.apache.comet.CometConf$. A session that registers CometSparkSessionExtensions through spark.sql.extensions without the plugin hit the same failure on the driver.

The relationship between the two configs still matters: a JVM shuffle batch larger than spark.comet.batchSize hands the native operators after the shuffle larger batches than configured. So the rule now lives where the values are read, instead of in a validator that runs during class initialization.

I chose to cap the value rather than fail at read time. Failing would still break the job from the issue: a user who only lowers spark.comet.batchSize to 4096 leaves spark.comet.shuffle.jvm.batchSize at its default of 8192, so every JVM shuffle write would fail. Capping keeps that job working and still guarantees what the check was for. It also means an explicit spark.comet.shuffle.jvm.batchSize larger than spark.comet.batchSize is now capped, where before it failed the shuffle write with an IllegalArgumentException. The config's doc string and the tuning guide say so.

What changes are included in this PR?

Stacked on #6287. Only the commit(s) after 54676eb are new. Related to #6287, which adds the check that the same value is positive, next to the validator removed here.

  • CometConf: remove the cross-config validator from COMET_SHUFFLE_JVM_BATCH_SIZE (the positivity check from fix: validate spark.comet.shuffle.jvm.batchSize and spark.comet.exec.memoryPool #6287 stays) and add jvmShuffleBatchSize(), which returns spark.comet.shuffle.jvm.batchSize capped at spark.comet.batchSize, both read from the same SQLConf. The doc string now describes the cap.
  • CometDiskBlockWriter and SpillWriter: the two places that read the JVM shuffle batch size call jvmShuffleBatchSize() instead, a one-line change in each. CometDiskBlockWriter uses the value for its row-count spill trigger, and SpillWriter.doSpilling passes it to writeSortedFileNative, which covers both the bypass-merge-sort and the sort-based JVM shuffle writers.
  • Docs: the memory tuning guide said the value "must not exceed" spark.comet.batchSize; it now says the value is capped. The spill trigger in the JVM shuffle contributor guide mentions the cap.

I audited the other entries in CometConf for the same hazard and found none. No other validator or default reads another config while CometConf initializes: the remaining validators only look at their own value, the ShimCometConf values are constants, and createWithEnvVarOrDefault reads only environment variables.

How are these changes tested?

  • CometConfSuite, "CometConf initializes when spark.comet.batchSize is below the JVM shuffle batch size": loads CometConf through a class loader that defines its own copy of every org.apache.comet class, so that CometConf's initializer runs again, under SQLConf.withExistingConf with spark.comet.batchSize=4096. This reproduces the class initialization failure inside the normal test JVM. Without the fix, loading it throws the error from the issue: ExceptionInInitializerError, caused by IllegalArgumentException: '8192' in spark.comet.shuffle.jvm.batchSize is invalid. ScalaTest aborts the whole run on ExceptionInInitializerError, so the test reports it as an ordinary failure instead.
  • CometConfSuite, "JVM shuffle batch size is capped at spark.comet.batchSize where it is read": the default and a larger explicit value are capped at 4096, a smaller value is used as it is, and the entry itself still returns its default of 8192.
  • DisableAQECometShuffleSuite, "JVM shuffle writes batches no larger than spark.comet.batchSize": runs a JVM columnar shuffle of 20000 rows (asserting that the plan has a JVM CometShuffleExchangeExec) with spark.comet.batchSize=4096 and the JVM shuffle batch size left at its default, and checks that no batch coming out of the shuffle has more than 4096 rows. With the writers still reading COMET_SHUFFLE_JVM_BATCH_SIZE directly, it failed with 8192 was not less than or equal to 4096.

I ran all of CometConfSuite plus the new DisableAQECometShuffleSuite test on the default profile (Spark 4.1, Scala 2.13, JDK 17): 21 tests, all passed. I did not run the end-to-end reproductions from the issue (a local-cluster executor, and a driver-only local[1] session without the plugin), because both need a JVM that has not loaded CometConf yet. The class loader test exercises the same initializer path.

…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
…here it is read

The entry for spark.comet.shuffle.jvm.batchSize had a validator that
compared the value with COMET_BATCH_SIZE.get(). createWithDefault runs
the validators on the default inside CometConf's static initializer, so
the check read SQLConf.get at whatever moment CometConf was loaded. On
an executor that is inside the first task, where SQLConf.get holds the
session's confs, so with spark.comet.batchSize below 8192 the default
failed the check, CometConf threw ExceptionInInitializerError, and every
later task on that executor failed with NoClassDefFoundError. A session
that registered CometSparkSessionExtensions without the plugin hit the
same failure on the driver.

The validator is removed, and the JVM shuffle writers now read the batch
size through CometConf.jvmShuffleBatchSize, which caps
spark.comet.shuffle.jvm.batchSize at spark.comet.batchSize. Capping
rather than failing keeps a job that only lowers spark.comet.batchSize
working, and it keeps what the check was for: the JVM shuffle never
writes batches larger than spark.comet.batchSize. No other CometConf
entry reads another config while CometConf initializes.

Closes apache#6286

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @andygrove lgtm

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Lowering spark.comet.batchSize below 8192 could break CometConf initialization. The prerequisite also addresses nonpositive shuffle batch sizes and invalid memory-pool names.
  • Design approach: Apply the batch-size cap when writers read configuration. Keep individual-value validation in the existing config builders.
  • Correctness / compatibility analysis: Both JVM shuffle writers reach the capped value. Checked Spark’s configuration propagation against sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Accepted memory-pool names match native support. No introduced P1/P2 issues found within this review.
  • Key design decisions: Both batch sizes come from the same SQLConf. The small helper centralizes the cap. Additional work occurs at writer construction and spill boundaries, with no added per-row processing.
  • Implementation sketch: Update both writer call sites, retain positivity validation, normalize pool names with Locale.ROOT, and add configuration and shuffle regression coverage. Documentation describes the cap.
  • Behavioral changes worth calling out: Explicitly oversized JVM shuffle batches are capped instead of rejected. Lower global batch sizes no longer poison class initialization. Invalid values fail when read, and supported mixed-case pool names are accepted.
  • Suggested improvements: None meeting the P1/P2 reporting threshold.

Reviewed full SHA 4d72e385cc5296d4de9d8dee3e87595942cec93f against base 22a07067bda272ed68adc67733cbaa7462a13c85. Covered all seven files and both commits in the full PR diff, including prerequisite 54676eb52afa793f4fb23bcc3463aae1fbf05907. The PR remains non-draft. Existing review, comments, threads and prerequisite discussion revealed no unresolved substantiated P1/P2 concerns.

Routed skills: review-comet-pr, review-comet-shuffle-pr, and review-comet-memory-pr.

Exact-head CI: 24 successful checks and 14 skipped, with no failures or pending checks. Successful checks include Linux Spark 4.1 suites, Rust tests and TPC-H/TPC-DS verification. The shuffle log confirms 513 passing tests, including the new batch-size regression. CI tested merge commit f353fcc56db1f206d66922f0b528ef149f35dd80, combining this head with the requested base. Spark SQL, Iceberg and macOS suites were skipped.

Validation: Isolated compilation of exact-head configuration sources against Spark 4.1.3, Scala 2.13.17 and JDK 21 passed all 20 CometConfSuite tests. A fresh-JVM probe initialized successfully at batch size 4096. Restoring the old validator reproduced ExceptionInInitializerError. This used a test-only copy of Utils.stringToSeq. No full local Maven/native build or end-to-end shuffle run was performed, and no performance benchmark was run. Project files remain unchanged.

@andygrove
andygrove added this pull request to the merge queue Sep 29, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:shuffle Shuffle (JVM and native) bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

spark.comet.batchSize below 8192 makes CometConf fail to initialize on every executor

3 participants