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 c314574ffe3 MINOR: Move BrokerReconfigurable to server module (#22409)
c314574ffe3 is described below

commit c314574ffe376926bdd0c95b33b12691d9ca750a
Author: Eric Chang <[email protected]>
AuthorDate: Mon Jun 29 16:01:58 2026 +0800

    MINOR: Move BrokerReconfigurable to server module (#22409)
    
    Move `BrokerReconfigurable` from `server-common` to the `server` module
    so it can use `AbstractKafkaConfig` directly.
    
    For non-server components, add a server-side named adapter instead of
    making the component depend on the broker reconfiguration interface.
    This PR adds `DynamicLogCleanerConfig` to bridge `LogCleaner` with the
    broker dynamic config pipeline.
    
    This keeps `storage` independent from the broker reconfiguration
    contract and follows the existing named-adapter style used by other
    dynamic broker config handlers.
    
    Ref: https://github.com/apache/kafka/pull/22353#discussion_r3307835672
    
    ### Testing
    
    - `./gradlew :server-common:compileJava :storage:compileJava
    :server:compileJava :core:compileScala :storage:test --tests
    org.apache.kafka.storage.internals.log.LogCleanerTest --tests
    org.apache.kafka.storage.internals.log.LogCleanerIntegrationTest`
    - `./gradlew :server:checkstyleMain :core:checkstyleMain
    :storage:checkstyleMain`
    
    Reviewers: Ken Huang <[email protected]>, Chia-Ping Tsai
     <[email protected]>
---
 .../scala/kafka/server/DynamicBrokerConfig.scala   |  7 ++-
 .../server/metadata/BrokerMetadataPublisher.scala  |  5 ++-
 .../kafka/server}/config/BrokerReconfigurable.java | 12 +++---
 .../server/config/DynamicLogCleanerConfig.java     | 50 ++++++++++++++++++++++
 .../config/DynamicProducerStateManagerConfig.java  |  6 +--
 .../kafka/storage/internals/log/LogCleaner.java    | 22 +++-------
 .../internals/log/LogCleanerIntegrationTest.java   | 18 ++++----
 .../storage/internals/log/LogCleanerTest.java      | 14 +++---
 8 files changed, 87 insertions(+), 47 deletions(-)

diff --git a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala 
b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala
index 224ebe48d0d..b407c2e55ba 100755
--- a/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala
+++ b/core/src/main/scala/kafka/server/DynamicBrokerConfig.scala
@@ -35,12 +35,11 @@ import org.apache.kafka.common.utils.internals.LogContext
 import org.apache.kafka.common.utils.internals.BufferSupplier
 import org.apache.kafka.common.utils.Utils
 import org.apache.kafka.common.utils.internals.ConfigUtils
-import org.apache.kafka.config
 import org.apache.kafka.network.SocketServer
 import org.apache.kafka.raft.KafkaRaftClient
 import org.apache.kafka.server.{DynamicThreadPool, ProcessRole}
 import org.apache.kafka.server.common.{ApiMessageAndVersion, 
DirectoryEventHandler}
-import org.apache.kafka.server.config.{DynamicConfig, 
DynamicProducerStateManagerConfig, ServerConfigs, ServerLogConfigs, 
DynamicBrokerConfig => JDynamicBrokerConfig}
+import org.apache.kafka.server.config.{BrokerReconfigurable => 
JBrokerReconfigurable, DynamicConfig, DynamicProducerStateManagerConfig, 
ServerConfigs, ServerLogConfigs, DynamicBrokerConfig => JDynamicBrokerConfig}
 import org.apache.kafka.server.log.remote.storage.RemoteLogManagerConfig
 import org.apache.kafka.server.metrics.{ClientTelemetryExporterPlugin, 
MetricConfigs}
 import org.apache.kafka.server.telemetry.{ClientTelemetry, 
ClientTelemetryExporterProvider}
@@ -231,7 +230,7 @@ class DynamicBrokerConfig(private val kafkaConfig: 
KafkaConfig) extends Logging
     reconfigurables.add(reconfigurable)
   }
 
-  def addBrokerReconfigurable(reconfigurable: config.BrokerReconfigurable): 
Unit = {
+  def addBrokerReconfigurable(reconfigurable: JBrokerReconfigurable): Unit = {
     verifyReconfigurableConfigs(reconfigurable.reconfigurableConfigs)
     brokerReconfigurables.add(new BrokerReconfigurable {
       override def reconfigurableConfigs: util.Set[String] = 
reconfigurable.reconfigurableConfigs
@@ -527,7 +526,7 @@ class DynamicBrokerConfig(private val kafkaConfig: 
KafkaConfig) extends Logging
 }
 
 /**
- * Implement [[config.BrokerReconfigurable]] instead.
+ * Implement [[JBrokerReconfigurable]] instead.
  */
 trait BrokerReconfigurable {
 
diff --git 
a/core/src/main/scala/kafka/server/metadata/BrokerMetadataPublisher.scala 
b/core/src/main/scala/kafka/server/metadata/BrokerMetadataPublisher.scala
index 472ccdf0f37..52cdb0dd4a0 100644
--- a/core/src/main/scala/kafka/server/metadata/BrokerMetadataPublisher.scala
+++ b/core/src/main/scala/kafka/server/metadata/BrokerMetadataPublisher.scala
@@ -35,6 +35,7 @@ import org.apache.kafka.metadata.KRaftMetadataCache
 import org.apache.kafka.metadata.publisher.{AclPublisher, 
DelegationTokenPublisher, DynamicClientQuotaPublisher, 
DynamicTopicClusterQuotaPublisher, ScramPublisher}
 import org.apache.kafka.server.common.MetadataVersion.MINIMUM_VERSION
 import org.apache.kafka.server.common.{FinalizedFeatures, ShareVersion}
+import org.apache.kafka.server.config.DynamicLogCleanerConfig
 import org.apache.kafka.server.fault.FaultHandler
 import org.apache.kafka.storage.internals.log.{UnifiedLog, LogManager => 
JLogManager}
 
@@ -341,7 +342,9 @@ class BrokerMetadataPublisher(
       // Make the LogCleaner available for reconfiguration. We can't do this 
prior to this
       // point because LogManager#startup creates the LogCleaner object, if
       // log.cleaner.enable is true. TODO: improve this (see KAFKA-13610)
-      
Option(logManager.cleaner).foreach(config.dynamicConfig.addBrokerReconfigurable)
+      Option(logManager.cleaner).foreach(cleaner =>
+        config.dynamicConfig.addBrokerReconfigurable(new 
DynamicLogCleanerConfig(cleaner))
+      )
     } catch {
       case t: Throwable => fatalFaultHandler.handleFault("Error starting 
LogManager", t)
     }
diff --git 
a/server-common/src/main/java/org/apache/kafka/config/BrokerReconfigurable.java 
b/server/src/main/java/org/apache/kafka/server/config/BrokerReconfigurable.java
similarity index 87%
rename from 
server-common/src/main/java/org/apache/kafka/config/BrokerReconfigurable.java
rename to 
server/src/main/java/org/apache/kafka/server/config/BrokerReconfigurable.java
index f17590317bf..4b31b60e0b2 100644
--- 
a/server-common/src/main/java/org/apache/kafka/config/BrokerReconfigurable.java
+++ 
b/server/src/main/java/org/apache/kafka/server/config/BrokerReconfigurable.java
@@ -14,9 +14,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
-package org.apache.kafka.config;
-
-import org.apache.kafka.common.config.AbstractConfig;
+package org.apache.kafka.server.config;
 
 import java.util.Set;
 
@@ -29,8 +27,8 @@ import java.util.Set;
  * The reconfiguration process follows three steps:
  * <ol>
  *   <li>Determining which configurations can be dynamically updated via 
{@link #reconfigurableConfigs()}</li>
- *   <li>Validating the new configuration before applying it via {@link 
#validateReconfiguration(AbstractConfig)}</li>
- *   <li>Applying the new configuration via {@link 
#reconfigure(AbstractConfig, AbstractConfig)}</li>
+ *   <li>Validating the new configuration before applying it via {@link 
#validateReconfiguration(AbstractKafkaConfig)}</li>
+ *   <li>Applying the new configuration via {@link 
#reconfigure(AbstractKafkaConfig, AbstractKafkaConfig)}</li>
  * </ol>
  */
 public interface BrokerReconfigurable {
@@ -53,7 +51,7 @@ public interface BrokerReconfigurable {
      *
      * @param newConfig the new configuration to validate
      */
-    void validateReconfiguration(AbstractConfig newConfig);
+    void validateReconfiguration(AbstractKafkaConfig newConfig);
 
     /**
      * Applies the new configuration.
@@ -63,5 +61,5 @@ public interface BrokerReconfigurable {
      * @param oldConfig the previous configuration
      * @param newConfig the new configuration to apply
      */
-    void reconfigure(AbstractConfig oldConfig, AbstractConfig newConfig);
+    void reconfigure(AbstractKafkaConfig oldConfig, AbstractKafkaConfig 
newConfig);
 }
diff --git 
a/server/src/main/java/org/apache/kafka/server/config/DynamicLogCleanerConfig.java
 
b/server/src/main/java/org/apache/kafka/server/config/DynamicLogCleanerConfig.java
new file mode 100644
index 00000000000..03090e38389
--- /dev/null
+++ 
b/server/src/main/java/org/apache/kafka/server/config/DynamicLogCleanerConfig.java
@@ -0,0 +1,50 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.server.config;
+
+import org.apache.kafka.storage.internals.log.CleanerConfig;
+import org.apache.kafka.storage.internals.log.LogCleaner;
+
+import java.util.Set;
+
+public class DynamicLogCleanerConfig implements BrokerReconfigurable {
+    private final LogCleaner logCleaner;
+
+    public DynamicLogCleanerConfig(LogCleaner logCleaner) {
+        this.logCleaner = logCleaner;
+    }
+
+    @Override
+    public Set<String> reconfigurableConfigs() {
+        return logCleaner.reconfigurableConfigs();
+    }
+
+    @Override
+    public void validateReconfiguration(AbstractKafkaConfig newConfig) {
+        logCleaner.validateReconfiguration(new CleanerConfig(newConfig));
+    }
+
+    @Override
+    public void reconfigure(AbstractKafkaConfig oldConfig, AbstractKafkaConfig 
newConfig) {
+        logCleaner.reconfigure(new CleanerConfig(oldConfig), new 
CleanerConfig(newConfig));
+    }
+
+    @Override
+    public String toString() {
+        return "DynamicLogCleanerConfig";
+    }
+}
diff --git 
a/server/src/main/java/org/apache/kafka/server/config/DynamicProducerStateManagerConfig.java
 
b/server/src/main/java/org/apache/kafka/server/config/DynamicProducerStateManagerConfig.java
index de9289e09a7..3a71c73e2c4 100644
--- 
a/server/src/main/java/org/apache/kafka/server/config/DynamicProducerStateManagerConfig.java
+++ 
b/server/src/main/java/org/apache/kafka/server/config/DynamicProducerStateManagerConfig.java
@@ -16,9 +16,7 @@
  */
 package org.apache.kafka.server.config;
 
-import org.apache.kafka.common.config.AbstractConfig;
 import org.apache.kafka.common.config.ConfigException;
-import org.apache.kafka.config.BrokerReconfigurable;
 import org.apache.kafka.coordinator.transaction.TransactionLogConfig;
 import org.apache.kafka.storage.internals.log.ProducerStateManagerConfig;
 
@@ -41,7 +39,7 @@ public class DynamicProducerStateManagerConfig implements 
BrokerReconfigurable {
     }
 
     @Override
-    public void validateReconfiguration(AbstractConfig newConfig) {
+    public void validateReconfiguration(AbstractKafkaConfig newConfig) {
         TransactionLogConfig transactionLogConfig = new 
TransactionLogConfig(newConfig);
         if (transactionLogConfig.producerIdExpirationMs() < 0)
             throw new 
ConfigException(TransactionLogConfig.PRODUCER_ID_EXPIRATION_MS_CONFIG + "cannot 
be less than 0, current value is " +
@@ -49,7 +47,7 @@ public class DynamicProducerStateManagerConfig implements 
BrokerReconfigurable {
     }
 
     @Override
-    public void reconfigure(AbstractConfig oldConfig, AbstractConfig 
newConfig) {
+    public void reconfigure(AbstractKafkaConfig oldConfig, AbstractKafkaConfig 
newConfig) {
         TransactionLogConfig transactionLogConfig = new 
TransactionLogConfig(newConfig);
         if (producerStateManagerConfig.producerIdExpirationMs() != 
transactionLogConfig.producerIdExpirationMs()) {
             log.info("Reconfigure {} from {} to {}",
diff --git 
a/storage/src/main/java/org/apache/kafka/storage/internals/log/LogCleaner.java 
b/storage/src/main/java/org/apache/kafka/storage/internals/log/LogCleaner.java
index 5b962f8fe45..577119db173 100644
--- 
a/storage/src/main/java/org/apache/kafka/storage/internals/log/LogCleaner.java
+++ 
b/storage/src/main/java/org/apache/kafka/storage/internals/log/LogCleaner.java
@@ -17,12 +17,10 @@
 package org.apache.kafka.storage.internals.log;
 
 import org.apache.kafka.common.TopicPartition;
-import org.apache.kafka.common.config.AbstractConfig;
 import org.apache.kafka.common.config.ConfigException;
 import org.apache.kafka.common.errors.KafkaStorageException;
 import org.apache.kafka.common.utils.Time;
 import org.apache.kafka.common.utils.internals.LogContext;
-import org.apache.kafka.config.BrokerReconfigurable;
 import org.apache.kafka.server.config.ServerConfigs;
 import org.apache.kafka.server.metrics.KafkaMetricsGroup;
 import org.apache.kafka.server.util.ShutdownableThread;
@@ -95,7 +93,7 @@ import java.util.stream.IntStream;
  *    tombstone deletion.</li>
  * </ol>
  */
-public class LogCleaner implements BrokerReconfigurable {
+public class LogCleaner {
     private static final Logger LOG = 
LoggerFactory.getLogger(LogCleaner.class);
 
     public static final Set<String> RECONFIGURABLE_CONFIGS = Set.of(
@@ -249,10 +247,6 @@ public class LogCleaner implements BrokerReconfigurable {
         cleanerManager.removeMetrics();
     }
 
-    /**
-     * @return A set of configs that is reconfigurable in LogCleaner
-     */
-    @Override
     public Set<String> reconfigurableConfigs() {
         return RECONFIGURABLE_CONFIGS;
     }
@@ -260,11 +254,10 @@ public class LogCleaner implements BrokerReconfigurable {
     /**
      * Validate the new cleaner threads num is reasonable.
      *
-     * @param newConfig A submitted new AbstractConfig instance that contains 
new cleaner config
+     * @param newConfig the submitted cleaner config
      */
-    @Override
-    public void validateReconfiguration(AbstractConfig newConfig) {
-        int numThreads = new CleanerConfig(newConfig).numThreads;
+    public void validateReconfiguration(CleanerConfig newConfig) {
+        int numThreads = newConfig.numThreads;
         int currentThreads = config.numThreads;
         if (numThreads < 1)
             throw new ConfigException("Log cleaner threads should be at least 
1");
@@ -286,12 +279,11 @@ public class LogCleaner implements BrokerReconfigurable {
      * @param oldConfig the old log cleaner config
      * @param newConfig the new log cleaner config reconfigured
      */
-    @Override
-    public void reconfigure(AbstractConfig oldConfig, AbstractConfig 
newConfig) {
-        config = new CleanerConfig(newConfig);
+    public void reconfigure(CleanerConfig oldConfig, CleanerConfig newConfig) {
+        config = newConfig;
 
         double maxIoBytesPerSecond = config.maxIoBytesPerSecond;
-        if (maxIoBytesPerSecond != 
oldConfig.getDouble(CleanerConfig.LOG_CLEANER_IO_MAX_BYTES_PER_SECOND_PROP)) {
+        if (maxIoBytesPerSecond != oldConfig.maxIoBytesPerSecond) {
             LOG.info("Updating logCleanerIoMaxBytesPerSecond: {}", 
maxIoBytesPerSecond);
             throttler.updateDesiredRatePerSec(maxIoBytesPerSecond);
         }
diff --git 
a/storage/src/test/java/org/apache/kafka/storage/internals/log/LogCleanerIntegrationTest.java
 
b/storage/src/test/java/org/apache/kafka/storage/internals/log/LogCleanerIntegrationTest.java
index d549fb42927..2adb17e139e 100644
--- 
a/storage/src/test/java/org/apache/kafka/storage/internals/log/LogCleanerIntegrationTest.java
+++ 
b/storage/src/test/java/org/apache/kafka/storage/internals/log/LogCleanerIntegrationTest.java
@@ -578,14 +578,14 @@ public class LogCleanerIntegrationTest {
         
oldConfigMap.put(CleanerConfig.LOG_CLEANER_IO_MAX_BYTES_PER_SECOND_PROP, 
currentConfig.maxIoBytesPerSecond);
         oldConfigMap.put(CleanerConfig.LOG_CLEANER_BACKOFF_MS_PROP, 
currentConfig.backoffMs);
         oldConfigMap.put(CleanerConfig.LOG_CLEANER_ENABLE_PROP, 
currentConfig.enableCleaner);
-        AbstractConfig oldAbstractConfig = new AbstractConfig(configDef, 
oldConfigMap);
+        CleanerConfig oldCleanerConfig = new CleanerConfig(new 
AbstractConfig(configDef, oldConfigMap));
 
         Map<String, Object> newConfigMap = new HashMap<>(oldConfigMap);
         newConfigMap.put(CleanerConfig.LOG_CLEANER_THREADS_PROP, 2);
         newConfigMap.put(CleanerConfig.LOG_CLEANER_IO_BUFFER_SIZE_PROP, 
100000);
-        AbstractConfig newAbstractConfig = new AbstractConfig(configDef, 
newConfigMap);
+        CleanerConfig newCleanerConfig = new CleanerConfig(new 
AbstractConfig(configDef, newConfigMap));
 
-        cleaner.reconfigure(oldAbstractConfig, newAbstractConfig);
+        cleaner.reconfigure(oldCleanerConfig, newCleanerConfig);
 
         assertEquals(2, cleaner.cleanerCount());
         checkLastCleaned("log", 0, firstDirty);
@@ -599,8 +599,8 @@ public class LogCleanerIntegrationTest {
         cleaner = makeCleaner(TOPIC_PARTITIONS, CLEANER_BACKOFF_MS, 
MIN_COMPACTION_LAG, SEGMENT_SIZE);
         cleaner.startup();
 
-        AbstractConfig config1Thread = makeReconfigureConfig(1);
-        AbstractConfig config2Thread = makeReconfigureConfig(2);
+        CleanerConfig config1Thread = makeReconfigureConfig(1);
+        CleanerConfig config2Thread = makeReconfigureConfig(2);
 
         var checkError = CompletableFuture.runAsync(() -> {
             var endtime = System.currentTimeMillis() + 
Duration.ofSeconds(5).toMillis();
@@ -614,8 +614,8 @@ public class LogCleanerIntegrationTest {
             var useOne = true;
             var endtime = System.currentTimeMillis() + 
Duration.ofSeconds(5).toMillis();
             while (System.currentTimeMillis() < endtime) {
-                AbstractConfig oldCfg = useOne ? config2Thread : config1Thread;
-                AbstractConfig newCfg = useOne ? config1Thread : config2Thread;
+                CleanerConfig oldCfg = useOne ? config2Thread : config1Thread;
+                CleanerConfig newCfg = useOne ? config1Thread : config2Thread;
                 cleaner.reconfigure(oldCfg, newCfg);
                 useOne = !useOne;
             }
@@ -625,7 +625,7 @@ public class LogCleanerIntegrationTest {
         updateCleaner.join();
     }
 
-    private AbstractConfig makeReconfigureConfig(int numThreads) {
+    private CleanerConfig makeReconfigureConfig(int numThreads) {
         // Extend CleanerConfig.CONFIG_DEF with message.max.bytes, which 
CleanerConfig(AbstractConfig)
         // reads via ServerConfigs.MESSAGE_MAX_BYTES_CONFIG but which is not 
part of CleanerConfig's own ConfigDef.
         ConfigDef configDef = new ConfigDef(CleanerConfig.CONFIG_DEF)
@@ -633,7 +633,7 @@ public class LogCleanerIntegrationTest {
                         DEFAULT_MAX_MESSAGE_SIZE, ConfigDef.Importance.MEDIUM, 
"");
         Map<String, Object> props = new HashMap<>();
         props.put(CleanerConfig.LOG_CLEANER_THREADS_PROP, numThreads);
-        return new AbstractConfig(configDef, props);
+        return new CleanerConfig(new AbstractConfig(configDef, props));
     }
 
     private void checkLastCleaned(String topic, int partitionId, long 
firstDirty) throws InterruptedException {
diff --git 
a/storage/src/test/java/org/apache/kafka/storage/internals/log/LogCleanerTest.java
 
b/storage/src/test/java/org/apache/kafka/storage/internals/log/LogCleanerTest.java
index e7dac1f2ffc..05e1761f960 100644
--- 
a/storage/src/test/java/org/apache/kafka/storage/internals/log/LogCleanerTest.java
+++ 
b/storage/src/test/java/org/apache/kafka/storage/internals/log/LogCleanerTest.java
@@ -2280,11 +2280,11 @@ public class LogCleanerTest {
     public void testReconfigureLogCleanerIoMaxBytesPerSecond() {
         var oldConfig = 
makeReconfigureConfig(Map.of(CleanerConfig.LOG_CLEANER_IO_MAX_BYTES_PER_SECOND_PROP,
 10000000D));
 
-        LogCleaner logCleaner = new LogCleaner(new CleanerConfig(oldConfig),
-            List.of(TestUtils.tempDirectory()),
-            new ConcurrentHashMap<>(),
-            new LogDirFailureChannel(1),
-            time) {
+        LogCleaner logCleaner = new LogCleaner(oldConfig,
+                List.of(TestUtils.tempDirectory()),
+                new ConcurrentHashMap<>(),
+                new LogDirFailureChannel(1),
+                time) {
             // shutdown() and startup() are called in LogCleaner.reconfigure().
             // Empty startup() and shutdown() to ensure that no unnecessary 
log cleaner threads remain after this test.
             @Override
@@ -2703,11 +2703,11 @@ public class LogCleanerTest {
         return cleaner.doClean(logToClean, currentTime + tombstoneRetentionMs 
+ 1).getKey();
     }
 
-    private AbstractConfig makeReconfigureConfig(Map<String, Object> 
overrides) {
+    private CleanerConfig makeReconfigureConfig(Map<String, Object> overrides) 
{
         ConfigDef configDef = new ConfigDef(CleanerConfig.CONFIG_DEF)
             .define(ServerConfigs.MESSAGE_MAX_BYTES_CONFIG, ConfigDef.Type.INT,
                 ServerLogConfigs.MAX_MESSAGE_BYTES_DEFAULT, 
ConfigDef.Importance.HIGH,
                 ServerConfigs.MESSAGE_MAX_BYTES_DOC);
-        return new AbstractConfig(configDef, new HashMap<>(overrides));
+        return new CleanerConfig(new AbstractConfig(configDef, new 
HashMap<>(overrides)));
     }
 }

Reply via email to