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]

Reply via email to