andygrove commented on code in PR #5634:
URL: https://github.com/apache/datafusion-comet/pull/5634#discussion_r4166815133
##########
spark/src/main/scala/org/apache/spark/Plugins.scala:
##########
@@ -104,7 +104,9 @@ object CometDriverPlugin extends Logging {
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 (conf.getBoolean(
+ CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key,
+ CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.defaultValue.get)) {
Review Comment:
This landed on its own as #6360. The plugin requires `spark.comet.enabled`
and `spark.comet.exec.enabled` through `getBooleanConf`, and `Comet plugin
installs its cache serializer only if Comet can scan the cache natively` covers
the three cases you listed: the key unset, Comet off, and execution off. So the
flip itself needs no change to the plugin.
#6537, now merged into this branch, applies the same rule to two more
startup settings. One is Comet shuffle enabled without one of Comet's shuffle
managers, where Comet disables itself. The other is Kryo that requires
registration but has not registered Comet's cached batch, which is covered in
the upgrade guide thread.
##########
spark/src/main/scala/org/apache/comet/CometConf.scala:
##########
@@ -278,7 +278,7 @@ object CometConf extends ShimCometConf {
"SparkContext, otherwise caching fails as soon as a block is
serialized, including " +
"the disk half of the default MEMORY_AND_DISK storage level.")
.booleanConf
- .createWithDefault(false)
+ .createWithDefault(true)
Review Comment:
Both are on this branch now.
Option 1 is #6538, merged here. `CometExecRule` already recorded a reason
when Spark's scan reads Comet's format because the native cache scan is off at
runtime or the cached plan records `Dataset.observe` metrics. It returned
before looking at the plan when Comet or its native execution was off, which is
the case of a session that turns Comet off after caching. It now records a
reason on each such scan, naming the cause and pointing to
`spark.comet.exec.inMemoryCache.enabled=false` at startup as the way to keep
Spark's format.
For the benchmark, `runAdaptiveBenchmark` in `CometInMemoryCacheBenchmark`
runs with Comet and AQE on and Comet's other settings at their defaults, over
the same 5M-row relation cached in Spark's format and in Comet's. Your first
case is a Spark operator above the native scan, reading through a
columnar-to-row transition. Here that operator is the aggregate with Comet's
turned off, standing in for any operator Comet does not support. On an AMD
Ryzen 9 7950X3D:
| Read | Spark's format | Comet's format | Relative |
| --- | --: | --: | --: |
| Row count only (0 of 6) | 39 ms | 16 ms | 2.4x |
| 1 of 6 columns | 50 ms | 27 ms | 1.8x |
| 3 of 6 columns | 97 ms | 112 ms | 0.9x |
| 6 of 6 columns | 303 ms | 299 ms | 1.0x |
So that path is not the slow one. The 3-of-6 read is three long columns, and
the 10% there is `zstd` decompression: with the `none` codec, Comet's format is
1.6x faster than Spark's on that read. With Comet operators above the scan, the
same four reads come out at 1.2x, 1.5x, 0.9x and 1.3x. The slow path is still
Spark's own cache scan reading Comet's format, 2.0x to 3.6x slower on this
machine, and that path now records a reason.
The run also turned up a correction to the guide. Under default settings,
Spark's cache scan does not feed Comet operators at all, because that bridge
needs `spark.comet.sparkToColumnar.enabled`. So with the feature off, the
operators directly above a cached relation run on Spark, and the published
native-scan numbers, which turn the bridge on, were not the feature-off
baseline. The guide now has both comparisons.
##########
docs/source/user-guide/latest/in-memory-cache.md:
##########
@@ -24,13 +24,13 @@ format that Comet operators read directly. Without it, a
cached table is stored
format and every scan of it has to convert each batch before Comet can
continue, which shows up in
the plan as a `CometSparkColumnarToColumnar` above the cache scan.
-This feature is **experimental and disabled by default**. Turn it on at
startup, alongside the rest
-of Comet's configuration:
+This feature is **experimental and enabled by default**. To turn it off, set
the config at startup,
+alongside the rest of Comet's configuration:
Review Comment:
The upgrade guide has a 1.2.0 entry now. It covers the format change, how to
keep Spark's format, and the cases where Comet keeps Spark's format without
being asked.
For Kryo, #6537 (merged into this branch) does more than document the
requirement. When Kryo requires registration and has not registered Comet's
cached batch, the plugin now keeps Spark's format rather than installing one
that Kryo would reject, and its startup warning says so. It asks a Kryo
instance built from the application's conf, so registrations made through
`CometKryoRegistrator`, another registrator or `spark.kryo.classesToRegister`
all count. Two new suites run this end to end with a `DISK_ONLY` cache. In one,
the application registers only Spark's cached batch. Without the change, its
cache fails with `Class is not registered:
org.apache.spark.sql.comet.execution.arrow.CometCachedBatch`, and with it the
cache is stored as `DefaultCachedBatch` and reads back correctly. In the other,
the application registers Comet's classes through
`spark.kryo.classesToRegister`, and the cache keeps Comet's format.
That also answers the legacy-key question for me. With the gate in place,
the flip changes no result and raises no new error under the same explicit
configuration. What it changes is which operators run natively and the storage
format underneath them. The [config
conventions](https://github.com/apache/datafusion-comet/blob/main/docs/source/contributor-guide/config_conventions.md#changing-the-behavior-of-an-existing-config)
exempt changes to which expressions and operators run natively, and
`spark.comet.exec.inMemoryCache.enabled=false` restores the old format exactly.
So the entry sits with the changes that need no legacy key, as the guide's
introduction allows, and it names that setting. The versioning policy does list
default changes among its examples, though. If you read that as applying even
when results don't change, I can add a `spark.comet.legacy.*` key, but it would
do exactly what `spark.comet.exec.inMemoryCache.enabled=false` already does.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]