This is an automated email from the ASF dual-hosted git repository.
chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 971c6b41310 MINOR: Avoid reflective class loading for share state
persister (#23001)
971c6b41310 is described below
commit 971c6b413103aa0db6e629fbe2eba32bd6024f52
Author: YANG-SYUAN CHOU <[email protected]>
AuthorDate: Fri Jul 31 20:45:04 2026 +0800
MINOR: Avoid reflective class loading for share state persister (#23001)
Remove unnecessary reflective class loading from
`createShareStatePersister` by comparing the configured class name
directly.
Also remove the corresponding `DefaultStatePersister` entry from
`reflect-config.json`.
Follow-up to #22900.
Reviewers: Andrew Schofield <[email protected]>, Chia-Ping Tsai
<[email protected]>, Sushant Mahajan <[email protected]>
---
core/src/main/scala/kafka/server/BrokerServer.scala | 8 ++++----
docker/native/native-image-configs/reflect-config.json | 4 ----
2 files changed, 4 insertions(+), 8 deletions(-)
diff --git a/core/src/main/scala/kafka/server/BrokerServer.scala
b/core/src/main/scala/kafka/server/BrokerServer.scala
index c73ca7e566a..88c11bd019a 100644
--- a/core/src/main/scala/kafka/server/BrokerServer.scala
+++ b/core/src/main/scala/kafka/server/BrokerServer.scala
@@ -751,16 +751,16 @@ class BrokerServer(
}
private def createShareStatePersister(): Persister = {
- if (config.shareGroupConfig.shareGroupPersisterClassName.nonEmpty) {
- val klass =
Utils.loadClass(config.shareGroupConfig.shareGroupPersisterClassName,
classOf[Object]).asInstanceOf[Class[Persister]]
- if (klass.getName.equals(classOf[DefaultStatePersister].getName)) {
+ val className = config.shareGroupConfig.shareGroupPersisterClassName
+ if (className.nonEmpty) {
+ if (className.equals(classOf[DefaultStatePersister].getName)) {
DefaultStatePersister.instance(
NetworkUtils.buildNetworkClient("Persister", config, metrics,
Time.SYSTEM, new LogContext(s"[Persister broker=${config.brokerId}]")),
new ShareCoordinatorMetadataCacheHelperImpl(metadataCache, key =>
shareCoordinator.partitionFor(key), config.interBrokerListenerName,
groupConfigManager, () => config.messageMaxBytes),
Time.SYSTEM,
shareGroupTimer
)
- } else if (klass.getName.equals(classOf[NoOpStatePersister].getName)) {
+ } else if (className.equals(classOf[NoOpStatePersister].getName)) {
info("Using no-op persister")
new NoOpStatePersister()
} else {
diff --git a/docker/native/native-image-configs/reflect-config.json
b/docker/native/native-image-configs/reflect-config.json
index 1dd72ae79b9..b36559f848b 100644
--- a/docker/native/native-image-configs/reflect-config.json
+++ b/docker/native/native-image-configs/reflect-config.json
@@ -1092,10 +1092,6 @@
"name":"org.apache.kafka.server.logger.LoggingControllerMBean",
"queryAllPublicMethods":true
},
-{
- "name":"org.apache.kafka.server.share.persister.DefaultStatePersister",
-
"methods":[{"name":"<init>","parameterTypes":["org.apache.kafka.server.share.persister.PersisterStateManager"]
}]
-},
{
"name":"org.apache.kafka.storage.internals.checkpoint.CleanShutdownFileHandler$Content",
"allDeclaredFields":true,