Skip to content

fix: keep Spark's cache format when Kryo would reject Comet's or Comet disables itself - #6537

Open
andygrove wants to merge 3 commits into
apache:mainfrom
andygrove:fix/cache-serializer-kryo-gate
Open

andygrove wants to merge 3 commits into
apache:mainfrom
andygrove:fix/cache-serializer-kryo-gate

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

No issue. Raised in the review of #5634, which turns the in-memory cache on by default: the strict-Kryo thread. Follows #6360.

Rationale for this change

#6360 made CometDriverPlugin install Comet's cache serializer only when Comet and its native execution are enabled at startup. spark.sql.cache.serializer is static, so an application that can never use Comet's format should keep Spark's. Two more startup settings decide whether Comet's format can work, and turning the cache on by default in #5634 would expose both:

  • Kryo with registration required. Under spark.kryo.registrationRequired=true, Kryo rejects any class it was not told about. Comet's cached batch is registered only through CometKryoRegistrator, and spark.kryo.registrator is read before any plugin runs, so Comet cannot add it. Without it, the first cached block Spark serializes (the disk half of the default MEMORY_AND_DISK, the _SER levels, replication) fails with Class is not registered: org.apache.spark.sql.comet.execution.arrow.CometCachedBatch, where Spark's own format works. The plugin only warned. Once the cache is on by default, that is a new error under an existing configuration.
  • Comet shuffle without Comet's shuffle manager. Comet disables itself when spark.comet.shuffle.enabled, which defaults to true, is on and spark.shuffle.manager is not one of Comet's managers. Such an application would cache everything in Comet's format with only Spark operators to read it, which is the case fix: install Comet's cache serializer only when Comet and native execution are enabled #6360 set out to exclude.

What changes are included in this PR?

  • CometDriverPlugin.maybeSetCacheSerializer also requires that Kryo can store Comet's cached batch, and, while Comet shuffle is enabled, that spark.shuffle.manager names CometShuffleManager or CometCelebornShuffleManager. Kryo can store it unless it requires registration and spark.kryo.registrator does not list CometKryoRegistrator. Otherwise caches keep Spark's format, as they do with the cache config off.
  • The plugin's startup warning for the Kryo case now says that Comet keeps Spark's cache format. Native broadcast still needs the registrator, as before.
  • The plugin's boolean config reads also check a deprecated alternative key, as a session's reads do, so spark.comet.exec.shuffle.enabled=false counts. This applies to the existing executor memory overhead warning too.
  • The in-memory cache guide, the installation guide's Kryo section, the plugin overview and the config's doc string describe the new conditions. The cache guide also notes that Spark registers its own cached batch with Kryo only from 4.1, and that CometKryoRegistrator registers it too.

How are these changes tested?

  • CometInMemoryCacheSuite: a new test covers the shuffle manager conditions (unset, sort, both of Comet's managers, Comet shuffle off under either key) and the Kryo conditions (no registrator, another registrator only, Comet's listed beside another, Kryo without registration required). The existing plugin tests now set Comet's shuffle manager, which they had relied on implicitly.
  • The new CometInMemoryCacheKryoUnregisteredSuite runs the setup from the review end to end: the plugin, Kryo with registrationRequired=true, and a registrator that registers Spark's cached batch but nothing of Comet's, as an application whose caches already work under registration on Spark 3.4 to 4.0 has to. A DISK_ONLY cache is stored as DefaultCachedBatch and reads back correctly. Without the change, the same test fails with Class is not registered: org.apache.spark.sql.comet.execution.arrow.CometCachedBatch. The suite is added to both PR workflows.
  • Ran locally on Spark 4.1: the plugin tests in CometInMemoryCacheSuite, CometInMemoryCacheKryoSuite, CometInMemoryCacheKryoUnregisteredSuite, the CometPlugins* suites and CometConfSuite. On Spark 3.5: the same cache suites and CometPluginsMemoryOverheadWarningSuite, plus the scalafix check.

…t disables itself

CometDriverPlugin installed ArrowCachedBatchSerializer whenever Comet, its
native execution and the cache config were enabled at startup. Two more
startup settings decide whether that format can work:

- With spark.kryo.registrationRequired=true and no CometKryoRegistrator,
  Kryo rejects CometCachedBatch, so caching failed with "Class is not
  registered" as soon as Spark serialized a cached block, where Spark's
  own format works. The plugin only warned. It now keeps Spark's format,
  and the warning says so.
- With Comet shuffle enabled but neither of Comet's shuffle managers
  configured, Comet disables itself, so every cache was stored in Comet's
  format with only Spark operators to read it. The plugin now keeps
  Spark's format there too.

The plugin's boolean config reads also honor a deprecated alternative
key, such as spark.comet.exec.shuffle.enabled, as a session does.
@github-actions github-actions Bot added the bug Something isn't working label Oct 2, 2026
The sentence described a session started with the default, which keeps
Spark's format only while the default is off.

@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: The plugin could install Comet’s cache format when strict Kryo rejected its batches or Comet disabled itself because of the shuffle manager.
  • Design approach: Add startup checks before installing the cache serializer and honor deprecated boolean configuration keys.
  • Correctness / compatibility analysis: The shuffle check and configuration precedence match existing behavior. One P2 regression remains: valid application-provided Kryo registrations can now trigger a switch to an unregistered Spark cache format on Spark 3.4–4.0.
  • Key design decisions: Serializer selection remains at startup, existing custom serializers remain untouched, and the helpers add no per-batch work.
  • Implementation sketch: Update the plugin, configuration documentation and cache guide. Add startup assertions and a disk-cache fallback suite registered in both PR workflows.
  • Behavioral changes worth calling out: Absence of the specific CometKryoRegistrator name now changes cache format even when the application has registered Comet’s classes another way.
  • Suggested improvements: Account for effective Kryo registrations before switching formats, and add a regression case for spark.kryo.classesToRegister.

Reviewed full head 73326f86925c1c5ae092a664c56d6f46285edc2a against base 1b488623a167d9c0ef325e2e48590ed7d81d8c46. The PR is not a draft. Reviewed all PR changes and the additional direct-comparison metrics differences, which originate from base-only commit #6314. Routed skills: review-comet-pr and review-comet-ffi-pr. There are no existing reviews, comments or threads on this PR.

Validation: Compared relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. A source-extracted Spark 3.5.9 probe reproduced the registration failure, and eight startup assertions passed. These probes used isolated scaffolding, not a full Comet/native build. Full repository suites were not run locally.

Exact-head CI at inspection: 17 successful, seven running, 14 skipped, no failures. Spark SQL, Iceberg and macOS checks were skipped. CI is not yet a completed green verdict.

Recommend request changes for the reproduced P2 below.

getBooleanConf(conf, CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED)) {
getBooleanConf(conf, CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED) &&
(!getBooleanConf(conf, CometConf.COMET_SHUFFLE_ENABLED) || isCometShuffleManager(conf)) &&
!isKryoRegistratorMissing(conf)) {

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.

[P2] Could this gate account for registrations supplied through spark.kryo.classesToRegister or an application’s own registrator? With Comet’s cache enabled and strict Kryo, an application can already register CometCachedBatch and its members without listing CometKryoRegistrator. Those batches serialize successfully before this change. The new gate instead selects DefaultCachedBatch, which Spark 3.4–4.0 does not register automatically. If the application registered only Comet’s format, a previously working DISK_ONLY cache now fails with Class is not registered: org.apache.spark.sql.execution.columnar.DefaultCachedBatch. Preserve usable registrations when selecting the format, and add coverage for this configuration.

Evidence: Reproduced on Spark 3.5.9 using the exact-head CometCachedBatch definition and Kryo predicate extracted into /tmp/comet6537-73326-review-e4lkum2_/KryoGateProbe.scala. With registrationRequired=true and Comet’s payload/member classes supplied through classesToRegister, output was head gate rejects=true, CometCachedBatch DISK_ONLY storage and reread passed, then Spark DISK_ONLY dataframe cache failed: unregistered DefaultCachedBatch. Separately compiled base/head startup methods selected ArrowCachedBatchSerializer before the change and DefaultCachedBatchSerializer afterward under the same registration configuration. Spark’s KryoSerializer explicitly supports classesToRegister and custom registrators, while built-in DefaultCachedBatch registration begins in 4.1.

The gate kept Spark's cache format whenever spark.kryo.registrator did
not list CometKryoRegistrator, even where the application had registered
Comet's cached batch another way, such as spark.kryo.classesToRegister.
Comet's format works there, while Spark's does not before 4.1, which
registers DefaultCachedBatch itself only from then on, so the gate broke
caches that worked.

Ask a Kryo instance built from the application's conf which of Comet's
classes it has registered, and keep Spark's format only if Comet's
cached batch is not among them. If Kryo cannot be built, fall back to
looking for CometKryoRegistrator in spark.kryo.registrator. The startup
warning uses the same check and names the missing classes.

@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: The plugin could select Comet’s cache format when strict Kryo rejected it or the shuffle configuration disabled Comet.
  • Design approach: Gate automatic serializer selection on startup configuration and inspect Kryo’s effective class registrations.
  • Correctness / compatibility analysis: The previous P2 concerning application-provided registrations is addressed. spark.kryo.classesToRegister and custom registrators now preserve Comet’s format. Relevant Spark sources were checked across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0.
  • Key design decisions: Registration inspection stays at startup, adding no per-batch work. The small helpers preserve custom serializer choices and give primary configuration keys precedence over deprecated alternatives.
  • Implementation sketch: Update plugin gates, registration diagnostics and documentation. Add startup assertions and disk-cache suites registered in both PR workflows.
  • Behavioral changes worth calling out: Ineligible applications retain Spark’s cache format. Applications with usable Comet registrations retain Comet’s format, even without naming CometKryoRegistrator.
  • Suggested improvements: None meeting the P1/P2 reporting threshold.

No introduced P1/P2 issues found within this review. No substantiated existing blockers remain.

Reviewed the full diff from base 1b488623a167d9c0ef325e2e48590ed7d81d8c46 to head 667ea36256fa8ef662f1b586a04fd9fcfa863ddd, including all 12 changed files and existing discussion. The PR remains non-draft. The metrics-file differences were inspected and traced to base-only commit #6314.

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

Exact-head CI: 25 successful checks, 15 skipped, no failures or running checks. Linux Spark 4.1 logs confirm both new cache suites passed. Spark SQL, Iceberg and macOS suites were skipped.

Validation: A Spark 3.5.9 probe using extracted head methods passed 20 startup assertions, Comet payload disk-storage/reread checks and Spark fallback dataframe-cache checks. This used isolated scaffolding, not a full local Comet/native build. Full repository suites and performance benchmarks were not run locally.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants