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)));
}
}