From f124ea4e162e58e031bd1dd669223b3f46edfc61 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Mon, 28 Sep 2026 16:38:17 -0600 Subject: [PATCH] fix: install Comet's cache serializer only when Comet and native execution 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. --- .../user-guide/latest/in-memory-cache.md | 5 ++- .../scala/org/apache/comet/CometConf.scala | 4 ++- .../main/scala/org/apache/spark/Plugins.scala | 9 ++++-- .../comet/exec/CometInMemoryCacheSuite.scala | 31 +++++++++++++++++++ 4 files changed, 45 insertions(+), 4 deletions(-) diff --git a/docs/source/user-guide/latest/in-memory-cache.md b/docs/source/user-guide/latest/in-memory-cache.md index 8e39f00fe5e..5704365ac9a 100644 --- a/docs/source/user-guide/latest/in-memory-cache.md +++ b/docs/source/user-guide/latest/in-memory-cache.md @@ -35,7 +35,10 @@ $SPARK_HOME/bin/spark-shell \ It has to be set before the `SparkContext` starts. Comet's driver plugin chooses `spark.sql.cache.serializer` once, while the context is initializing, so a session that started -with the default goes on using Spark's cache format however the config is set afterwards. +with the default goes on using Spark's cache format however the config is set afterwards. The +plugin installs Comet's serializer only if `spark.comet.enabled` and `spark.comet.exec.enabled` +are enabled at that point too, because an application that starts without native execution could +not scan Comet's format natively. ## What changes when it is enabled diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index ae971883e62..7e9e75d4564 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -265,7 +265,9 @@ object CometConf extends ShimCometConf { .category(CATEGORY_EXEC) .doc("Whether to enable Comet native execution for in-memory cached tables. Its value at " + "startup also decides whether CometDriverPlugin installs Comet's cache serializer, " + - "which stores cached data in Arrow format. Because spark.sql.cache.serializer is a " + + "which stores cached data in Arrow format. The plugin installs it only if " + + "spark.comet.enabled and spark.comet.exec.enabled are also enabled at startup. " + + "Because spark.sql.cache.serializer is a " + "static config, the cached format is fixed for the application, and disabling this " + "at runtime only sends cached scans back to Spark's execution path. Relations whose " + "schema Comet's Arrow writer does not support are always cached in Spark's default " + diff --git a/spark/src/main/scala/org/apache/spark/Plugins.scala b/spark/src/main/scala/org/apache/spark/Plugins.scala index b0680ae695d..6d54faf2a6b 100644 --- a/spark/src/main/scala/org/apache/spark/Plugins.scala +++ b/spark/src/main/scala/org/apache/spark/Plugins.scala @@ -102,13 +102,18 @@ object CometDriverPlugin extends Logging { /** Spark config key under which the loaded Comet version is exposed at runtime. */ val COMET_VERSION_CONFIG = "spark.comet.version" - // Use Comet's cache serializer only for the native in-memory cache path. + // Use Comet's cache serializer only when the native in-memory cache scan can run, which needs + // Comet and its native execution as well as the cache config. spark.sql.cache.serializer is + // static, so an application that starts with Comet or native execution off would otherwise + // store every cache in Comet's format, with only Spark operators to read it. // If the application already set spark.sql.cache.serializer, leave that value // unchanged so Comet does not replace a user-selected cache format. private[apache] def maybeSetCacheSerializer( conf: SparkConf, extraConfs: ju.HashMap[String, String]): Unit = { - if (conf.getBoolean(CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key, false)) { + if (getBooleanConf(conf, CometConf.COMET_ENABLED) && + getBooleanConf(conf, CometConf.COMET_EXEC_ENABLED) && + getBooleanConf(conf, CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED)) { val serializerKey = StaticSQLConf.SPARK_CACHE_SERIALIZER.key val serializerValue = "org.apache.spark.sql.comet.execution.arrow.ArrowCachedBatchSerializer" diff --git a/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala index b12ef02c5f7..58ca444cd23 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala @@ -984,6 +984,37 @@ class CometInMemoryCacheSuite extends CometTestBase { assert(!userExtraConfs.containsKey(serializerKey)) } + test("Comet plugin installs its cache serializer only if Comet can scan the cache natively") { + val serializerKey = StaticSQLConf.SPARK_CACHE_SERIALIZER.key + + def installed(settings: (String, String)*): Boolean = { + val conf = new SparkConf().setAll(settings) + val extraConfs = new ju.HashMap[String, String]() + CometDriverPlugin.maybeSetCacheSerializer(conf, extraConfs) + assert(conf.contains(serializerKey) == extraConfs.containsKey(serializerKey)) + extraConfs.containsKey(serializerKey) + } + + val cometOn = CometConf.COMET_ENABLED.key -> "true" + val execOn = CometConf.COMET_EXEC_ENABLED.key -> "true" + val cacheOn = CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "true" + + assert(installed(cometOn, execOn, cacheOn)) + // An application that starts with Comet or its native execution off can never plan + // CometInMemoryTableScan, and spark.sql.cache.serializer is static, so its caches keep + // Spark's format. + assert(!installed(CometConf.COMET_ENABLED.key -> "false", execOn, cacheOn)) + assert(!installed(cometOn, CometConf.COMET_EXEC_ENABLED.key -> "false", cacheOn)) + // Unset keys take their defaults rather than values of the plugin's own. + assert( + installed(cometOn, execOn) == + CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.defaultValue.get) + assert( + installed(cacheOn) == + (CometConf.COMET_ENABLED.defaultValue.get && + CometConf.COMET_EXEC_ENABLED.defaultValue.get)) + } + test("Comet in-memory cache supports empty projection scan") { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",