andygrove opened a new pull request, #6366: URL: https://github.com/apache/datafusion-comet/pull/6366
## 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 54676eb52 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 #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. -- 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]
