Conversation
…onf init The check on spark.comet.shuffle.jvm.batchSize read spark.comet.batchSize while CometConf initialized. An executor first loads CometConf inside a task, so a batch size below 8192 failed the default and left CometConf unable to initialize there. Cap the JVM shuffle batch size where it is read instead, and test that no config check reads another config.
Member
Contributor
Author
|
correct, we made PRs at almost the same time |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Which issue does this PR close?
Closes #6286.
Rationale for this change
createWithDefaultruns an entry's checks on its default whileCometConfinitializes. The check onspark.comet.shuffle.jvm.batchSizecompared the value withCOMET_BATCH_SIZE.get(), which reads whateverSQLConf.getreturns at that moment. An executor first loadsCometConfinside a task, where that is the session's conf, so aspark.comet.batchSizebelow 8192 failed the default and leftCometConfunable to initialize for the rest of the executor's life. The driver hits the same failure when Comet is loaded only throughspark.sql.extensions.The existing tests missed it for two reasons:
CometConfis initialized once, on the driver, during suite setup, whileSQLConf.getstill returns the defaults. The tests that set a smaller batch size throughwithSQLConfnever run the initializer again, and no suite starts separate executor JVMs.ConfigEntryWithDefault.getreturns the default without running the checks. So those tests never compared the default with their smaller batch size, and JVM shuffle kept writing batches of up to 8192 rows under them. The limit the check was meant to enforce never applied to the default.What changes are included in this PR?
spark.comet.shuffle.jvm.batchSizeis removed. The newCometConf.jvmShuffleBatchSizereturns that value capped atspark.comet.batchSize.CometDiskBlockWriterandSpillWriter, the only two readers, call it.checkValuenow says that a check also runs on the default whileCometConfinitializes, so it must not read other configs.Capping the value where it is read, rather than failing there, is what lets the default work with a smaller
spark.comet.batchSize. Two things behave differently as a result:spark.comet.shuffle.jvm.batchSizelarger thanspark.comet.batchSizeis now capped. Since 0.14.0 it failed the task.spark.comet.batchSizebelow 8192 already worked, as in local mode, JVM shuffle now writes batches of at most that many rows instead of 8192.This overlaps with #6287, which adds a positivity check to the same entry on the line above the one removed here. Whichever PR merges second should keep both.
How are these changes tested?
Two new tests in
CometConfSuite:JVM shuffle batch size is capped at spark.comet.batchSizecovers the default, the issue's case of a 4096 batch size with the default JVM shuffle batch size, and explicit JVM shuffle batch sizes above and below the batch size.config checks do not read other configsruns every registered entry's checks on its default, ascreateWithDefaultdoes, underSQLConf.withExistingConfwith a conf that fails the test on anygetConfString. That is the rule spark.comet.batchSize below 8192 makes CometConf fail to initialize on every executor #6286 broke, and the test applies it to every entry, not only this one. The check this PR removes readspark.comet.batchSizethrough that conf.There is no end-to-end test. The executor case needs
CometConfto be loaded for the first time inside a task, which takes separate executor JVMs, as in the issue'slocal-clusterreproduction. In Comet's local-mode suites,CometConfis already loaded by the time a test sets a smaller batch size, so neither the executor case nor the driver case can be reproduced there.