sunchao commented on code in PR #6537:
URL: https://github.com/apache/datafusion-comet/pull/6537#discussion_r4165226887
##########
spark/src/main/scala/org/apache/spark/Plugins.scala:
##########
@@ -174,15 +175,22 @@ object CometDriverPlugin extends Logging {
// 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.
+ // store every cache in Comet's format, with only Spark operators to read
it. So would one that
+ // leaves Comet shuffle enabled without Comet's shuffle manager, since Comet
then disables
+ // itself.
+ // Nor is it used where Kryo would reject its cached batches, under
+ // spark.kryo.registrationRequired without CometKryoRegistrator: caching
that works in Spark's
+ // format would then fail the first time Spark serialized a cached block.
// 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 (getBooleanConf(conf, CometConf.COMET_ENABLED) &&
getBooleanConf(conf, CometConf.COMET_EXEC_ENABLED) &&
- 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)) {
Review Comment:
[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.
--
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]