Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion docs/source/user-guide/latest/in-memory-cache.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
4 changes: 3 additions & 1 deletion spark/src/main/scala/org/apache/comet/CometConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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 " +
Expand Down
9 changes: 7 additions & 2 deletions spark/src/main/scala/org/apache/spark/Plugins.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Loading