From a12f79c5d2f2795c0241737c5a4f709ca01dc48e Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Wed, 2 Sep 2026 16:37:09 +0800 Subject: [PATCH 1/6] [SPARK-59169][CORE] Log a warning when KVStoreProtobufSerializer falls back to JSON SerDe Assisted-by: Claude Fable 5 --- .../protobuf/KVStoreProtobufSerializer.scala | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala index d87c9e6d59a75..ea06ae2539c74 100644 --- a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala +++ b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala @@ -19,9 +19,12 @@ package org.apache.spark.status.protobuf import java.lang.reflect.ParameterizedType import java.util.ServiceLoader +import java.util.concurrent.ConcurrentHashMap import scala.jdk.CollectionConverters._ +import org.apache.spark.internal.Logging +import org.apache.spark.internal.LogKeys.CLASS_NAME import org.apache.spark.status.KVUtils.KVStoreScalaSerializer private[spark] class KVStoreProtobufSerializer extends KVStoreScalaSerializer { @@ -39,7 +42,7 @@ private[spark] class KVStoreProtobufSerializer extends KVStoreScalaSerializer { } } -private[spark] object KVStoreProtobufSerializer { +private[spark] object KVStoreProtobufSerializer extends Logging { private[this] lazy val serializerMap: Map[Class[_], ProtobufSerDe[Any]] = { def getGenericsType(klass: Class[_]): Class[_] = { @@ -51,6 +54,14 @@ private[spark] object KVStoreProtobufSerializer { }.toMap } - def getSerializer(klass: Class[_]): Option[ProtobufSerDe[Any]] = - serializerMap.get(klass) + private[this] val missedClasses = ConcurrentHashMap.newKeySet[Class[_]]() + + def getSerializer(klass: Class[_]): Option[ProtobufSerDe[Any]] = { + val serializer = serializerMap.get(klass) + if (serializer.isEmpty && missedClasses.add(klass)) { + logWarning(log"No Protobuf SerDe found for class ${MDC(CLASS_NAME, klass.getName)}, " + + log"falling back to use the JSON SerDe.") + } + serializer + } } From 6976b6836b811761195152e4fa635181c52dc135 Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Wed, 2 Sep 2026 19:49:00 +0800 Subject: [PATCH 2/6] nit: remove log prefix --- .../spark/status/protobuf/KVStoreProtobufSerializer.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala index ea06ae2539c74..f9f28442a4173 100644 --- a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala +++ b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala @@ -60,7 +60,7 @@ private[spark] object KVStoreProtobufSerializer extends Logging { val serializer = serializerMap.get(klass) if (serializer.isEmpty && missedClasses.add(klass)) { logWarning(log"No Protobuf SerDe found for class ${MDC(CLASS_NAME, klass.getName)}, " + - log"falling back to use the JSON SerDe.") + "falling back to use the JSON SerDe.") } serializer } From 7dd4abe0e20d3330cd77d9443f93f29f1f67371a Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Wed, 2 Sep 2026 20:45:15 +0800 Subject: [PATCH 3/6] address comment, add UT Assisted-by: Claude Fable 5 --- .../protobuf/KVStoreProtobufSerializer.scala | 2 +- .../KVStoreProtobufSerializerSuite.scala | 18 ++++++++++++++++++ 2 files changed, 19 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala index f9f28442a4173..ea06ae2539c74 100644 --- a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala +++ b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala @@ -60,7 +60,7 @@ private[spark] object KVStoreProtobufSerializer extends Logging { val serializer = serializerMap.get(klass) if (serializer.isEmpty && missedClasses.add(klass)) { logWarning(log"No Protobuf SerDe found for class ${MDC(CLASS_NAME, klass.getName)}, " + - "falling back to use the JSON SerDe.") + log"falling back to use the JSON SerDe.") } serializer } diff --git a/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala b/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala index 23cb99bce2ff4..be7a682499c48 100644 --- a/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala +++ b/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala @@ -22,6 +22,8 @@ import java.util.Date import scala.collection.mutable import scala.io.Source +import org.apache.logging.log4j.Level + import org.apache.spark.{JobExecutionStatus, SparkFunSuite} import org.apache.spark.executor.ExecutorMetrics import org.apache.spark.metrics.ExecutorMetricType @@ -35,6 +37,20 @@ import org.apache.spark.util.Utils.tryWithResource class KVStoreProtobufSerializerSuite extends SparkFunSuite { private val serializer = new KVStoreProtobufSerializer() + test("SPARK-59169: log a warning once per class when no ProtobufSerDe is found") { + val appender = new LogAppender("KVStoreProtobufSerializer fallback warning") + withLogAppender(appender, loggerNames = Seq(classOf[KVStoreProtobufSerializer].getName)) { + serializer.serialize(FallbackTestData("a")) + serializer.serialize(FallbackTestData("b")) + } + val warnings = appender.loggingEvents + .filter(_.getLevel == Level.WARN) + .map(_.getMessage.getFormattedMessage) + .filter(_.contains(classOf[FallbackTestData].getName)) + assert(warnings.size === 1) + assert(warnings.head.contains("No Protobuf SerDe found for class")) + } + test("All the string fields must be optional to avoid NPE") { val protoFile = getWorkspaceFilePath( "core", "src", "main", "protobuf", "org", "apache", "spark", "status", "protobuf", @@ -1703,3 +1719,5 @@ class KVStoreProtobufSerializerSuite extends SparkFunSuite { } } } + +private[protobuf] case class FallbackTestData(value: String) From cea4510c2b915f33123a46a4e2d8170064ba64c7 Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Wed, 2 Sep 2026 22:22:13 +0800 Subject: [PATCH 4/6] make UT robust against test ordering Assisted-by: Claude Fable 5 --- .../spark/status/protobuf/KVStoreProtobufSerializer.scala | 2 ++ .../spark/status/protobuf/KVStoreProtobufSerializerSuite.scala | 1 + 2 files changed, 3 insertions(+) diff --git a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala index ea06ae2539c74..d9aa1e671beb0 100644 --- a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala +++ b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala @@ -56,6 +56,8 @@ private[spark] object KVStoreProtobufSerializer extends Logging { private[this] val missedClasses = ConcurrentHashMap.newKeySet[Class[_]]() + private[protobuf] def resetMissedClassesForTesting(): Unit = missedClasses.clear() + def getSerializer(klass: Class[_]): Option[ProtobufSerDe[Any]] = { val serializer = serializerMap.get(klass) if (serializer.isEmpty && missedClasses.add(klass)) { diff --git a/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala b/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala index be7a682499c48..6065246df5826 100644 --- a/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala +++ b/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala @@ -38,6 +38,7 @@ class KVStoreProtobufSerializerSuite extends SparkFunSuite { private val serializer = new KVStoreProtobufSerializer() test("SPARK-59169: log a warning once per class when no ProtobufSerDe is found") { + KVStoreProtobufSerializer.resetMissedClassesForTesting() val appender = new LogAppender("KVStoreProtobufSerializer fallback warning") withLogAppender(appender, loggerNames = Seq(classOf[KVStoreProtobufSerializer].getName)) { serializer.serialize(FallbackTestData("a")) From 9912008a19b18378886e6d4cd0955bf1c0384498 Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Wed, 2 Sep 2026 23:02:56 +0800 Subject: [PATCH 5/6] skip by-design JSON classes, fix wording Assisted-by: Claude Fable 5 --- .../protobuf/KVStoreProtobufSerializer.scala | 15 +++++++++++++-- .../protobuf/KVStoreProtobufSerializerSuite.scala | 10 ++++++++++ 2 files changed, 23 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala index d9aa1e671beb0..b4420bbb069e3 100644 --- a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala +++ b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala @@ -23,9 +23,12 @@ import java.util.concurrent.ConcurrentHashMap import scala.jdk.CollectionConverters._ +import org.apache.spark.deploy.history.FsHistoryProviderMetadata import org.apache.spark.internal.Logging import org.apache.spark.internal.LogKeys.CLASS_NAME +import org.apache.spark.status.AppStatusStoreMetadata import org.apache.spark.status.KVUtils.KVStoreScalaSerializer +import org.apache.spark.util.kvstore.{LevelDB, RocksDB} private[spark] class KVStoreProtobufSerializer extends KVStoreScalaSerializer { override def serialize(o: Object): Array[Byte] = @@ -56,13 +59,21 @@ private[spark] object KVStoreProtobufSerializer extends Logging { private[this] val missedClasses = ConcurrentHashMap.newKeySet[Class[_]]() + // KVStore bookkeeping values fall back to the JSON SerDe by design: they are tiny, + // written at most once per store open, and gain nothing from Protobuf. Skip warning. + private[this] val jsonByDesignClasses: Set[Class[_]] = Set( + classOf[AppStatusStoreMetadata], + classOf[FsHistoryProviderMetadata], + classOf[LevelDB.TypeAliases], + classOf[RocksDB.TypeAliases]) + private[protobuf] def resetMissedClassesForTesting(): Unit = missedClasses.clear() def getSerializer(klass: Class[_]): Option[ProtobufSerDe[Any]] = { val serializer = serializerMap.get(klass) - if (serializer.isEmpty && missedClasses.add(klass)) { + if (serializer.isEmpty && !jsonByDesignClasses.contains(klass) && missedClasses.add(klass)) { logWarning(log"No Protobuf SerDe found for class ${MDC(CLASS_NAME, klass.getName)}, " + - log"falling back to use the JSON SerDe.") + log"falling back to the JSON SerDe.") } serializer } diff --git a/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala b/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala index 6065246df5826..b0facd7ce80e7 100644 --- a/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala +++ b/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala @@ -52,6 +52,16 @@ class KVStoreProtobufSerializerSuite extends SparkFunSuite { assert(warnings.head.contains("No Protobuf SerDe found for class")) } + test("SPARK-59169: no warning for KVStore bookkeeping classes without ProtobufSerDe") { + KVStoreProtobufSerializer.resetMissedClassesForTesting() + val appender = new LogAppender("KVStoreProtobufSerializer by-design fallback") + withLogAppender(appender, loggerNames = Seq(classOf[KVStoreProtobufSerializer].getName)) { + serializer.serialize(AppStatusStoreMetadata(1L)) + } + val warnings = appender.loggingEvents.filter(_.getLevel == Level.WARN) + assert(warnings.isEmpty) + } + test("All the string fields must be optional to avoid NPE") { val protoFile = getWorkspaceFilePath( "core", "src", "main", "protobuf", "org", "apache", "spark", "status", "protobuf", From 2b73e608716808738984b5814bd4c51edf28dc67 Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Wed, 2 Sep 2026 23:38:58 +0800 Subject: [PATCH 6/6] keep warning for metadata classes, skip only kvstore internals Assisted-by: Claude Fable 5 --- .../status/protobuf/KVStoreProtobufSerializer.scala | 9 +++------ .../status/protobuf/KVStoreProtobufSerializerSuite.scala | 4 +++- 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala index b4420bbb069e3..d8bb13e316f6b 100644 --- a/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala +++ b/core/src/main/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializer.scala @@ -23,10 +23,8 @@ import java.util.concurrent.ConcurrentHashMap import scala.jdk.CollectionConverters._ -import org.apache.spark.deploy.history.FsHistoryProviderMetadata import org.apache.spark.internal.Logging import org.apache.spark.internal.LogKeys.CLASS_NAME -import org.apache.spark.status.AppStatusStoreMetadata import org.apache.spark.status.KVUtils.KVStoreScalaSerializer import org.apache.spark.util.kvstore.{LevelDB, RocksDB} @@ -59,11 +57,10 @@ private[spark] object KVStoreProtobufSerializer extends Logging { private[this] val missedClasses = ConcurrentHashMap.newKeySet[Class[_]]() - // KVStore bookkeeping values fall back to the JSON SerDe by design: they are tiny, - // written at most once per store open, and gain nothing from Protobuf. Skip warning. + // The KVStore backends' own bookkeeping values fall back to the JSON SerDe by design: + // they are internals of the kvstore library, which the ProtobufSerDe SPI does not cover. + // Skip warning for them. private[this] val jsonByDesignClasses: Set[Class[_]] = Set( - classOf[AppStatusStoreMetadata], - classOf[FsHistoryProviderMetadata], classOf[LevelDB.TypeAliases], classOf[RocksDB.TypeAliases]) diff --git a/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala b/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala index b0facd7ce80e7..ef983e1b2b6d1 100644 --- a/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala +++ b/core/src/test/scala/org/apache/spark/status/protobuf/KVStoreProtobufSerializerSuite.scala @@ -33,6 +33,7 @@ import org.apache.spark.status._ import org.apache.spark.status.api.v1._ import org.apache.spark.ui.scope.{RDDOperationEdge, RDDOperationNode} import org.apache.spark.util.Utils.tryWithResource +import org.apache.spark.util.kvstore.{LevelDB, RocksDB} class KVStoreProtobufSerializerSuite extends SparkFunSuite { private val serializer = new KVStoreProtobufSerializer() @@ -56,7 +57,8 @@ class KVStoreProtobufSerializerSuite extends SparkFunSuite { KVStoreProtobufSerializer.resetMissedClassesForTesting() val appender = new LogAppender("KVStoreProtobufSerializer by-design fallback") withLogAppender(appender, loggerNames = Seq(classOf[KVStoreProtobufSerializer].getName)) { - serializer.serialize(AppStatusStoreMetadata(1L)) + assert(KVStoreProtobufSerializer.getSerializer(classOf[RocksDB.TypeAliases]).isEmpty) + assert(KVStoreProtobufSerializer.getSerializer(classOf[LevelDB.TypeAliases]).isEmpty) } val warnings = appender.loggingEvents.filter(_.getLevel == Level.WARN) assert(warnings.isEmpty)