fix: install Comet's cache serializer only when Comet and native execution are enabled - #6360
Conversation
…ution are enabled CometDriverPlugin installed ArrowCachedBatchSerializer whenever spark.comet.exec.inMemoryCache.enabled was true at startup, including in applications that start with spark.comet.enabled=false or spark.comet.exec.enabled=false. spark.sql.cache.serializer is static, so those applications stored every cache in Comet's format without ever being able to scan it natively, and only Spark operators read it. Require both configs as well, reading all three through getBooleanConf so that an unset key takes its default rather than a hard-coded false.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: The driver plugin could install Comet’s cache serializer while Comet or native execution was disabled, forcing Spark readers to consume Comet’s cache format.
- Design approach: Require all three enable flags before automatically installing the serializer.
- Correctness / compatibility analysis: Checked Spark’s
StaticSQLConf,SharedState, plugin initialization, andInMemoryRelationsources across 3.4.3, 3.5.9, 4.0.4, 4.1.3, and experimental 4.2.0. The startup decision matches their static serializer semantics. Comet’s existing fallback handles Spark-format caches after execution is enabled later. - Key design decisions: Reuses
getBooleanConf, respects configuration defaults, and preserves custom serializers. The additional checks run during startup, without adding per-query or per-batch work or a new abstraction. - Implementation sketch: Adds the three-flag condition, regression assertions for enabled, disabled, and unset configurations, and updates configuration and user documentation.
- Behavioral changes worth calling out: Starting with Comet or native execution disabled retains Spark’s cache format even if execution is enabled later.
- Suggested improvements: No introduced P1/P2 issues found within this review.
Reviewed all four changed files in the full diff from 65b334bfdd25196091a42d773bc3e61d19212e52 to f124ea4e162e58e031bd1dd669223b3f46edfc61. The PR was open and not a draft. Both the supplied snapshot and live discussion checks contained no reviews, issue comments, or review threads.
Routed skills: review-comet-pr. No sibling skill applied to the changed scope.
Exact-head CI: 15 checks passed, including Spark 4.1 compilation and lint checks. Native build and Rust tests remained in progress. Thirteen checks were skipped, including Spark SQL, Iceberg, and macOS suites. No failures were reported at inspection time.
Validation limits: git diff --check passed. Local Scala/native suites were not run because this checkout has no built JVM/native artifacts. The successful CI compilation explicitly skipped tests, so runtime validation remains limited. No project files were changed and nothing was published.
Which issue does this PR close?
Part of #5485. It is the first of the changes promised in the #5634 review before the in-memory cache is turned on by default.
Rationale for this change
CometDriverPlugininstallsArrowCachedBatchSerializerasspark.sql.cache.serializerwheneverspark.comet.exec.inMemoryCache.enabledis true at startup, even when the application starts withspark.comet.enabled=falseorspark.comet.exec.enabled=false.spark.sql.cache.serializeris static, so every cache in such an application is stored in Comet's format but can never be scanned byCometInMemoryTableScan. Spark operators read all of it, which is the slow path described in #5485 and in the Limitations section of the in-memory cache guide.Keeping the plugin in
spark.pluginscluster-wide and switching Comet off per application is a common setup, and @mbutrovich pointed out on #5634 that flipping the default would give every such application Comet's format. #5485 lists this check as its second option.What changes are included in this PR?
maybeSetCacheSerializernow also requiresspark.comet.enabledandspark.comet.exec.enabled. All three keys are read through the existinggetBooleanConfhelper, so an unset key takes its config default. That includes the cache key itself, which was read with a hard-codedfalse, so the plugin follows the default when it changes.A session that starts with either config off and turns it on later keeps Spark's format.
CometExecRulealready records a fallback reason for a relation cached with another serializer.The config's doc string and the in-memory cache guide now say when the plugin installs the serializer.
How are these changes tested?
A new test in
CometInMemoryCacheSuitecallsmaybeSetCacheSerializerwith all three configs on (installed), with Comet off and with native execution off (not installed), and with the keys unset (each follows its config default). It also checks that the driver conf and theextraConfssent to executors agree. With the old condition restored, the Comet-off case fails.I ran the two plugin tests in
CometInMemoryCacheSuiteand all ofCometPluginsSuiteon Spark 4.1 with Scala 2.13 and on Spark 3.4 with Scala 2.12.