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]

Reply via email to