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 44e45841cc5 KAFKA-18628 Deprecate broker.id config (KIP-1232) (#22977)
44e45841cc5 is described below
commit 44e45841cc5a07a5354da99752dae37f4e4b48b6
Author: Ming-Yen Chung <[email protected]>
AuthorDate: Sat Aug 1 00:03:47 2026 +0800
KAFKA-18628 Deprecate broker.id config (KIP-1232) (#22977)
Implement [KIP-1232](https://cwiki.apache.org/confluence/x/Hgp3Fw).
- Deprecate the `broker.id` configuration for removal in Kafka 5.0.
- Log a deprecation warning at broker startup when `broker.id` is
explicitly set, pointing users to `node.id`.
- Pass `node.id` to the tiered storage plugins alongside `broker.id`,
and switch Kafka's own `RemoteStorageManager` and
`RemoteLogMetadataManager` implementations to read it. `broker.id` is
still passed for compatibility and will be dropped in 5.0.
- Reword the `node.id` / `broker.id` mismatch error so it points at
removing `broker.id` rather than keeping both in sync.
- Clean up redundant `broker.id` settings in test files — `node.id`
alone is enough via the existing synonym mechanism.
The startup warning is checked on the raw properties in
`KafkaConfig.fromProps`, because `populateSynonyms` copies `node.id`
into `broker.id`, so `originals()` cannot tell whether the user actually
set it. It is not guarded by `doLog`, since the broker startup path
calls `fromProps` with `doLog = false`.
The tiered storage plugin change is not covered by the KIP yet — I will
amend KIP-1232 to include it.
Test result:
```
# broker.id=1
❯ bin/kafka-server-start.sh /tmp/kip1232.properties
[2026-07-28 23:49:14,050] INFO Registered
`kafka:type=kafka.Log4jController` MBean
(org.apache.kafka.server.logger.Log4jControllerRegistration)
[2026-07-28 23:49:14,114] WARN The 'broker.id' configuration is
deprecated and will be removed in Apache Kafka 5.0. Please use 'node.id'
instead. (kafka.server.KafkaConfig$)
# node.id=1
❯ bin/kafka-server-start.sh config/server.properties
[2026-07-28 23:52:45,352] INFO Registered
`kafka:type=kafka.Log4jController` MBean
(org.apache.kafka.server.logger.Log4jControllerRegistration)
[2026-07-28 23:52:45,461] INFO Registered signal handlers for TERM, INT,
HUP (org.apache.kafka.common.utils.internals.LoggingSignalHandler)
[2026-07-28 23:52:45,462] INFO [ControllerServer id=1] Starting
controller (kafka.server.ControllerServer)
```
Reviewers: Gaurav Narula <[email protected]>, Ken Huang
<[email protected]>, Chia-Ping Tsai <[email protected]>, Federico
Valeri <[email protected]>
---
core/src/main/scala/kafka/server/KafkaConfig.scala | 14 ++++++--
.../api/AbstractAuthorizerIntegrationTest.scala | 1 -
.../kafka/server/DescribeClusterRequestTest.scala | 3 +-
.../scala/unit/kafka/server/KafkaConfigTest.scala | 37 +++++++++++++++++++++-
.../kafka/server/KafkaMetricsReporterTest.scala | 2 --
.../test/scala/unit/kafka/utils/TestUtils.scala | 1 -
docs/getting-started/upgrade.md | 2 ++
.../kafka/jmh/util/BenchmarkConfigUtils.java | 1 -
.../apache/kafka/server/config/ServerConfigs.java | 6 +++-
.../kafka/server/config/AbstractKafkaConfig.java | 2 ++
.../apache/kafka/api/EndToEndClusterIdTest.java | 4 +--
.../server/config/AbstractKafkaConfigTest.java | 2 ++
.../remote/storage/RemoteLogMetadataManager.java | 4 ++-
.../log/remote/storage/RemoteStorageManager.java | 4 ++-
.../TopicBasedRemoteLogMetadataManagerConfig.java | 27 +++++++++++++++-
.../log/remote/storage/RemoteLogManager.java | 11 +++++--
.../storage/RemoteLogMetadataManagerTestUtils.java | 4 +--
...picBasedRemoteLogMetadataManagerConfigTest.java | 33 +++++++++++++++++--
.../TopicBasedRemoteLogMetadataManagerTest.java | 2 +-
.../log/remote/storage/LocalTieredStorageTest.java | 2 +-
.../log/remote/storage/RemoteLogManagerTest.java | 9 ++++--
.../log/remote/storage/LocalTieredStorage.java | 32 +++++++++++++++----
22 files changed, 169 insertions(+), 34 deletions(-)
diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala
b/core/src/main/scala/kafka/server/KafkaConfig.scala
index 03f2c46a890..a80c9ea1590 100755
--- a/core/src/main/scala/kafka/server/KafkaConfig.scala
+++ b/core/src/main/scala/kafka/server/KafkaConfig.scala
@@ -49,7 +49,7 @@ import scala.jdk.CollectionConverters._
import scala.collection.{Map, Seq}
import scala.jdk.OptionConverters.{RichOption, RichOptional}
-object KafkaConfig {
+object KafkaConfig extends Logging {
def main(args: Array[String]): Unit = {
val combined = new ConfigDef(configDef)
@@ -73,8 +73,15 @@ object KafkaConfig {
def fromProps(props: Properties): KafkaConfig =
fromProps(props, true)
- def fromProps(props: Properties, doLog: Boolean): KafkaConfig =
+ def fromProps(props: Properties, doLog: Boolean): KafkaConfig = {
+ // Checked on the raw properties because populateSynonyms copies node.id
into broker.id.
+ // Do not add a doLog guard: the broker startup path calls this with doLog
= false.
+ if (props.containsKey(ServerConfigs.BROKER_ID_CONFIG)) {
+ warn(s"The '${ServerConfigs.BROKER_ID_CONFIG}' configuration is
deprecated and will be removed in " +
+ s"Apache Kafka 5.0. Please use '${KRaftConfigs.NODE_ID_CONFIG}'
instead.")
+ }
new KafkaConfig(props, doLog)
+ }
def fromProps(defaults: Properties, overrides: Properties): KafkaConfig =
fromProps(defaults, overrides, true)
@@ -427,7 +434,8 @@ class KafkaConfig private(doLog: Boolean, val props:
util.Map[_, _])
private def validateValues(): Unit = {
if (nodeId != brokerId) {
- throw new ConfigException(s"You must set
`${KRaftConfigs.NODE_ID_CONFIG}` to the same value as
`${ServerConfigs.BROKER_ID_CONFIG}`.")
+ throw new ConfigException(s"`${KRaftConfigs.NODE_ID_CONFIG}` and
`${ServerConfigs.BROKER_ID_CONFIG}` must be set to the same value. " +
+ s"`${ServerConfigs.BROKER_ID_CONFIG}` is deprecated, please use
`${KRaftConfigs.NODE_ID_CONFIG}` instead.")
}
require(logRollTimeMillis >= 1, "log.roll.ms must be greater than or equal
to 1")
require(logRollTimeJitterMillis >= 0, "log.roll.jitter.ms must be greater
than or equal to 0")
diff --git
a/core/src/test/scala/integration/kafka/api/AbstractAuthorizerIntegrationTest.scala
b/core/src/test/scala/integration/kafka/api/AbstractAuthorizerIntegrationTest.scala
index 0c510850b94..b8a447e7ee5 100644
---
a/core/src/test/scala/integration/kafka/api/AbstractAuthorizerIntegrationTest.scala
+++
b/core/src/test/scala/integration/kafka/api/AbstractAuthorizerIntegrationTest.scala
@@ -100,7 +100,6 @@ class AbstractAuthorizerIntegrationTest extends
BaseRequestTest {
consumerConfig.setProperty(ConsumerConfig.GROUP_ID_CONFIG, group)
override def brokerPropertyOverrides(properties: Properties): Unit = {
- properties.put(ServerConfigs.BROKER_ID_CONFIG, brokerId.toString)
addNodeProperties(properties)
}
diff --git
a/core/src/test/scala/unit/kafka/server/DescribeClusterRequestTest.scala
b/core/src/test/scala/unit/kafka/server/DescribeClusterRequestTest.scala
index 96283d60389..a24a5604ae9 100644
--- a/core/src/test/scala/unit/kafka/server/DescribeClusterRequestTest.scala
+++ b/core/src/test/scala/unit/kafka/server/DescribeClusterRequestTest.scala
@@ -24,6 +24,7 @@ import
org.apache.kafka.common.requests.{DescribeClusterRequest, DescribeCluster
import org.apache.kafka.common.resource.ResourceType
import org.apache.kafka.common.utils.Utils
import org.apache.kafka.coordinator.group.GroupCoordinatorConfig
+import org.apache.kafka.raft.KRaftConfigs
import org.apache.kafka.security.authorizer.AclEntry
import org.apache.kafka.server.config.{ServerConfigs, ReplicationConfigs}
import org.junit.jupiter.api.Assertions.{assertEquals, assertTrue}
@@ -38,7 +39,7 @@ class DescribeClusterRequestTest extends BaseRequestTest {
override def brokerPropertyOverrides(properties: Properties): Unit = {
properties.setProperty(GroupCoordinatorConfig.OFFSETS_TOPIC_PARTITIONS_CONFIG,
"1")
properties.setProperty(ReplicationConfigs.DEFAULT_REPLICATION_FACTOR_CONFIG,
"2")
- properties.setProperty(ServerConfigs.BROKER_RACK_CONFIG,
s"rack/${properties.getProperty(ServerConfigs.BROKER_ID_CONFIG)}")
+ properties.setProperty(ServerConfigs.BROKER_RACK_CONFIG,
s"rack/${properties.getProperty(KRaftConfigs.NODE_ID_CONFIG)}")
}
@BeforeEach
diff --git a/core/src/test/scala/unit/kafka/server/KafkaConfigTest.scala
b/core/src/test/scala/unit/kafka/server/KafkaConfigTest.scala
index 34f48635581..4fd07e4e9c5 100755
--- a/core/src/test/scala/unit/kafka/server/KafkaConfigTest.scala
+++ b/core/src/test/scala/unit/kafka/server/KafkaConfigTest.scala
@@ -45,6 +45,8 @@ import org.apache.logging.log4j.Level
import org.junit.jupiter.api.Assertions._
import org.junit.jupiter.api.Test
import org.junit.jupiter.api.function.Executable
+import org.junit.jupiter.params.ParameterizedTest
+import org.junit.jupiter.params.provider.ValueSource
import scala.jdk.CollectionConverters._
import scala.util.Using
@@ -1661,7 +1663,8 @@ class KafkaConfigTest {
props.setProperty(ServerConfigs.BROKER_ID_CONFIG, "1")
props.setProperty(KRaftConfigs.NODE_ID_CONFIG, "2")
props.setProperty(KRaftConfigs.CONTROLLER_LISTENER_NAMES_CONFIG,
"CONTROLLER")
- assertEquals("You must set `node.id` to the same value as `broker.id`.",
+ assertEquals("`node.id` and `broker.id` must be set to the same value. " +
+ "`broker.id` is deprecated, please use `node.id` instead.",
assertThrows(classOf[ConfigException], () =>
KafkaConfig.fromProps(props)).getMessage())
}
@@ -1710,6 +1713,38 @@ class KafkaConfigTest {
assertEquals("3", originals.get(KRaftConfigs.NODE_ID_CONFIG))
}
+ // The warning must fire regardless of doLog, since the broker startup path
passes doLog = false.
+ @ParameterizedTest
+ @ValueSource(booleans = Array(true, false))
+ def testBrokerIdDeprecationWarning(doLog: Boolean): Unit = {
+ val deprecationWarning = "The 'broker.id' configuration is deprecated and
will be removed in " +
+ "Apache Kafka 5.0. Please use 'node.id' instead."
+
+ // Register on the KafkaConfig logger rather than the root logger, so that
the warning is only
+ // captured if it is actually logged by `object KafkaConfig`.
+ Using.resource(LogCaptureAppender.createAndRegister(KafkaConfig.getClass))
{ appender =>
+ appender.setClassLogger(KafkaConfig.getClass, Level.WARN)
+ // The appender cannot be reset, so the counts asserted below are
cumulative.
+ def warningCount: Int = appender.getMessages.asScala.count(_ ==
deprecationWarning)
+
+ // Only node.id set: no warning.
+ val props = new Properties()
+ props.putAll(kraftProps())
+ KafkaConfig.fromProps(props, doLog)
+ assertEquals(0, warningCount)
+
+ // Both broker.id and node.id set to the same value: deprecation warning.
+ props.setProperty(ServerConfigs.BROKER_ID_CONFIG, "3")
+ KafkaConfig.fromProps(props, doLog)
+ assertEquals(1, warningCount)
+
+ // Only broker.id set: deprecation warning.
+ props.remove(KRaftConfigs.NODE_ID_CONFIG)
+ KafkaConfig.fromProps(props, doLog)
+ assertEquals(2, warningCount)
+ }
+ }
+
@Test
def testSaslJwksEndpointRetryDefaults(): Unit = {
val props = new Properties()
diff --git
a/core/src/test/scala/unit/kafka/server/KafkaMetricsReporterTest.scala
b/core/src/test/scala/unit/kafka/server/KafkaMetricsReporterTest.scala
index e4bcb3d49a5..0ab80f32628 100644
--- a/core/src/test/scala/unit/kafka/server/KafkaMetricsReporterTest.scala
+++ b/core/src/test/scala/unit/kafka/server/KafkaMetricsReporterTest.scala
@@ -21,7 +21,6 @@ import java.util.concurrent.atomic.AtomicReference
import kafka.utils.TestUtils
import org.apache.kafka.common.metrics.{KafkaMetric, MetricsContext,
MetricsReporter}
import org.apache.kafka.common.utils.Utils
-import org.apache.kafka.server.config.ServerConfigs
import org.apache.kafka.server.metrics.MetricConfigs
import org.apache.kafka.test.{TestUtils => JTestUtils}
import org.junit.jupiter.api.{AfterEach, BeforeEach, Test, TestInfo}
@@ -73,7 +72,6 @@ class KafkaMetricsReporterTest extends QuorumTestHarness {
super.setUp(testInfo)
val props = TestUtils.createBrokerConfig(1)
props.setProperty(MetricConfigs.METRIC_REPORTER_CLASSES_CONFIG,
"kafka.server.KafkaMetricsReporterTest$MockMetricsReporter")
- props.setProperty(ServerConfigs.BROKER_ID_CONFIG, "1")
config = KafkaConfig.fromProps(props)
broker = createBroker(config, threadNamePrefix =
Option(this.getClass.getName))
broker.startup()
diff --git a/core/src/test/scala/unit/kafka/utils/TestUtils.scala
b/core/src/test/scala/unit/kafka/utils/TestUtils.scala
index 1f3bc0de45d..7f61b8872d3 100755
--- a/core/src/test/scala/unit/kafka/utils/TestUtils.scala
+++ b/core/src/test/scala/unit/kafka/utils/TestUtils.scala
@@ -235,7 +235,6 @@ object TestUtils extends Logging {
props.put(ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG, "true")
props.setProperty(KRaftConfigs.SERVER_MAX_STARTUP_TIME_MS_CONFIG,
TimeUnit.MINUTES.toMillis(10).toString)
props.put(KRaftConfigs.NODE_ID_CONFIG, nodeId.toString)
- props.put(ServerConfigs.BROKER_ID_CONFIG, nodeId.toString)
props.put(SocketServerConfigs.ADVERTISED_LISTENERS_CONFIG, listeners)
props.put(SocketServerConfigs.LISTENERS_CONFIG, listeners)
props.put(KRaftConfigs.CONTROLLER_LISTENER_NAMES_CONFIG, "CONTROLLER")
diff --git a/docs/getting-started/upgrade.md b/docs/getting-started/upgrade.md
index 857e0397dcb..61d32837bf2 100644
--- a/docs/getting-started/upgrade.md
+++ b/docs/getting-started/upgrade.md
@@ -46,6 +46,8 @@ type: docs
* The broker-side OAUTHBEARER JWT validator now fails fast at startup when a
JWKS endpoint (`sasl.oauthbearer.jwks.endpoint.url`) is configured but
`sasl.oauthbearer.expected.audience` or `sasl.oauthbearer.expected.issuer` is
not set. Brokers that previously started without these settings will now fail
to start until they are configured. To intentionally accept tokens regardless
of their audience or issuer, set the new
`sasl.oauthbearer.allow.unverified.audience` or `sasl.oauthbearer.a [...]
* When clients connect to the cluster, they now include cluster and node
information to enable detection and handling of misrouted connections. For
further details, please refer to
[KIP-1242](https://cwiki.apache.org/confluence/x/W4LMFw).
* The `kafka-cluster.sh` tool now provides an `api-versions` command to
display the API versions supported by the brokers or controllers, and it
accepts both `--bootstrap-server` and `--bootstrap-controller`. As a result,
`kafka-broker-api-versions.sh` is deprecated and will be removed in the next
major release; use `kafka-cluster.sh api-versions` instead. For further
details, please refer to
[KIP-1220](https://cwiki.apache.org/confluence/x/-QkbFw).
+ * The `broker.id` configuration is deprecated and will be removed in Kafka
5.0. Please use `node.id` instead. For further details, please refer to
[KIP-1232](https://cwiki.apache.org/confluence/x/Hgp3Fw).
+ * Tiered storage plugins are now configured with `node.id` in addition to
`broker.id`. Since `broker.id` will no longer be passed to
`RemoteStorageManager` and `RemoteLogMetadataManager` implementations in Kafka
5.0, plugins should read `node.id` instead. For further details, please refer
to [KIP-1232](https://cwiki.apache.org/confluence/x/Hgp3Fw).
* Brokers can now record a human-readable description of each streams
group's processing topology via a pluggable backend, retrievable through
`Admin#describeStreamsGroups` and `kafka-streams-groups.sh --describe
--topology`. The feature is disabled unless the new broker configuration
`group.streams.topology.description.plugin.class` is set to a
`StreamsGroupTopologyDescriptionPlugin` implementation; on the client side, the
new Kafka Streams configuration `topology.description.push.ena [...]
* The `kafka-producer-perf-test.sh` tool now supports `--record-key-range`,
`--key-distribution`, and `--random-seed` options to control the distribution
of record keys. Use `--key-distribution range` for sequential key assignment
(round-robin over the key range) or `--key-distribution random` for random key
selection. The `--random-seed` option allows reproducible benchmark runs when
using random key distribution. For further details, please refer to
[KIP-1299](https://cwiki.apache.or [...]
* Share groups now support dead-letter queue functionality as outlined in
[KIP-1191](https://cwiki.apache.org/confluence/x/fApJFg). Any records which are
released (beyond max delivery count) or rejected by the share consumer become
eligible for DLQ. Share group DLQ gets enabled when the Kafka feature
`share.version` is upgraded to 2. The user can configure a DLQ topic on a share
group by setting the dynamic config `errors.deadletterqueue.topic.name`
(default `""`) to the name of the DL [...]
diff --git
a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/util/BenchmarkConfigUtils.java
b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/util/BenchmarkConfigUtils.java
index 1a9fbbc30e2..28b40790c82 100644
---
a/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/util/BenchmarkConfigUtils.java
+++
b/jmh-benchmarks/src/main/java/org/apache/kafka/jmh/util/BenchmarkConfigUtils.java
@@ -39,7 +39,6 @@ public class BenchmarkConfigUtils {
props.put(ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG, "true");
props.setProperty(KRaftConfigs.SERVER_MAX_STARTUP_TIME_MS_CONFIG,
String.valueOf(TimeUnit.MINUTES.toMillis(10)));
props.put(KRaftConfigs.NODE_ID_CONFIG, "0");
- props.put(ServerConfigs.BROKER_ID_CONFIG, "0");
props.put(SocketServerConfigs.ADVERTISED_LISTENERS_CONFIG,
"PLAINTEXT://localhost:9092");
props.put(SocketServerConfigs.LISTENERS_CONFIG,
"PLAINTEXT://localhost:9092,CONTROLLER://localhost:9093");
diff --git
a/server-common/src/main/java/org/apache/kafka/server/config/ServerConfigs.java
b/server-common/src/main/java/org/apache/kafka/server/config/ServerConfigs.java
index 04a2353591f..939feaf7248 100644
---
a/server-common/src/main/java/org/apache/kafka/server/config/ServerConfigs.java
+++
b/server-common/src/main/java/org/apache/kafka/server/config/ServerConfigs.java
@@ -39,9 +39,13 @@ import static
org.apache.kafka.common.config.ConfigDef.Type.STRING;
public class ServerConfigs {
/** ********* General Configuration ***********/
+ @Deprecated(since = "4.4", forRemoval = true)
public static final String BROKER_ID_CONFIG = "broker.id";
+ @Deprecated(since = "4.4", forRemoval = true)
public static final int BROKER_ID_DEFAULT = -1;
- public static final String BROKER_ID_DOC = "The broker id for this
server.";
+ @Deprecated(since = "4.4", forRemoval = true)
+ public static final String BROKER_ID_DOC = "The broker id for this server.
" +
+ "This configuration is deprecated and will be removed in Apache
Kafka 5.0. Please use <code>node.id</code> instead.";
public static final String MESSAGE_MAX_BYTES_CONFIG = "message.max.bytes";
public static final String MESSAGE_MAX_BYTES_DOC =
TopicConfig.MAX_MESSAGE_BYTES_DOC +
diff --git
a/server/src/main/java/org/apache/kafka/server/config/AbstractKafkaConfig.java
b/server/src/main/java/org/apache/kafka/server/config/AbstractKafkaConfig.java
index 6220651ed03..9329173d2c7 100644
---
a/server/src/main/java/org/apache/kafka/server/config/AbstractKafkaConfig.java
+++
b/server/src/main/java/org/apache/kafka/server/config/AbstractKafkaConfig.java
@@ -138,6 +138,7 @@ public abstract class AbstractKafkaConfig extends
AbstractConfig {
return getInt(ServerConfigs.BACKGROUND_THREADS_CONFIG);
}
+ @SuppressWarnings("removal") // broker.id is deprecated (KIP-1232), but
this method stays and will read node.id in 5.0
public int brokerId() {
return getInt(ServerConfigs.BROKER_ID_CONFIG);
}
@@ -220,6 +221,7 @@ public abstract class AbstractKafkaConfig extends
AbstractConfig {
/**
* Copy a configuration map, populating some keys that we want to treat as
synonyms.
*/
+ @SuppressWarnings("removal") // broker.id is deprecated (KIP-1232), but it
still works as another name for node.id until 5.0
public static Map<Object, Object> populateSynonyms(Map<?, ?> input) {
Map<Object, Object> output = new HashMap<>(input);
Object brokerId = output.get(ServerConfigs.BROKER_ID_CONFIG);
diff --git
a/server/src/test/java/org/apache/kafka/api/EndToEndClusterIdTest.java
b/server/src/test/java/org/apache/kafka/api/EndToEndClusterIdTest.java
index 6c48f138a3f..beb2bca8205 100644
--- a/server/src/test/java/org/apache/kafka/api/EndToEndClusterIdTest.java
+++ b/server/src/test/java/org/apache/kafka/api/EndToEndClusterIdTest.java
@@ -30,7 +30,7 @@ import org.apache.kafka.common.test.ClusterInstance;
import org.apache.kafka.common.test.api.ClusterConfigProperty;
import org.apache.kafka.common.test.api.ClusterTest;
import org.apache.kafka.common.test.api.ClusterTestDefaults;
-import org.apache.kafka.server.config.ServerConfigs;
+import org.apache.kafka.raft.KRaftConfigs;
import org.apache.kafka.server.metrics.MetricConfigs;
import org.apache.kafka.test.MockConsumerInterceptor;
import org.apache.kafka.test.MockDeserializer;
@@ -102,7 +102,7 @@ public class EndToEndClusterIdTest {
String roles = (String) configs.get("process.roles");
if (roles == null) return;
- String id = (String) configs.get(ServerConfigs.BROKER_ID_CONFIG);
+ String id = (String) configs.get(KRaftConfigs.NODE_ID_CONFIG);
controllerId = roles.contains("controller") ? id : null;
brokerId = roles.contains("broker") ? id : null;
}
diff --git
a/server/src/test/java/org/apache/kafka/server/config/AbstractKafkaConfigTest.java
b/server/src/test/java/org/apache/kafka/server/config/AbstractKafkaConfigTest.java
index fae45548fa8..6661d2db593 100644
---
a/server/src/test/java/org/apache/kafka/server/config/AbstractKafkaConfigTest.java
+++
b/server/src/test/java/org/apache/kafka/server/config/AbstractKafkaConfigTest.java
@@ -52,6 +52,7 @@ public class AbstractKafkaConfigTest {
assertEquals(Collections.emptyMap(),
AbstractKafkaConfig.populateSynonyms(Collections.emptyMap()));
}
+ @SuppressWarnings("removal") // this test sets broker.id, which is
deprecated (KIP-1232) but still works until 5.0
@Test
public void testPopulateSynonymsOnMapWithoutNodeId() {
Map<String, String> input = new HashMap<>();
@@ -62,6 +63,7 @@ public class AbstractKafkaConfigTest {
assertEquals(expectedOutput,
AbstractKafkaConfig.populateSynonyms(input));
}
+ @SuppressWarnings("removal") // this test sets broker.id, which is
deprecated (KIP-1232) but still works until 5.0
@Test
public void testPopulateSynonymsOnMapWithoutBrokerId() {
Map<String, String> input = new HashMap<>();
diff --git
a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java
b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java
index 4e30c59bbc0..a1456dbce03 100644
---
a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java
+++
b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java
@@ -46,8 +46,10 @@ import java.util.concurrent.CompletableFuture;
* <code>remote.log.metadata.manager.listener.name</code> property is about
listener name of the local broker to which
* it should get connected if needed by RemoteLogMetadataManager
implementation.
* </p>
- * "cluster.id", "broker.id" and all other properties prefixed with the
config: "remote.log.metadata.manager.impl.prefix"
+ * "cluster.id", "node.id" and all other properties prefixed with the config:
"remote.log.metadata.manager.impl.prefix"
* (default value is "rlmm.config.") are passed when {@link #configure(Map)}
is invoked on this instance.
+ * "broker.id" is also passed with the same value as "node.id", but it is
deprecated since 4.4 and will not be passed
+ * from Apache Kafka 5.0 onwards. Implementations should read "node.id"
instead.
* <p>
*
* Implement {@link org.apache.kafka.common.metrics.Monitorable} to enable the
manager to register metrics.
diff --git
a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteStorageManager.java
b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteStorageManager.java
index 91252e0b021..c5689b1ba67 100644
---
a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteStorageManager.java
+++
b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteStorageManager.java
@@ -38,8 +38,10 @@ import java.util.Optional;
* This allows {@link RemoteStorageManager} to have eventual consistency on
metadata (although the data is stored
* in strongly consistent semantics).
* <p>
- * All properties prefixed with the config:
"remote.log.storage.manager.impl.prefix"
+ * "node.id" and all properties prefixed with the config:
"remote.log.storage.manager.impl.prefix"
* (default value is "rsm.config.") are passed when {@link #configure(Map)} is
invoked on this instance.
+ * "broker.id" is also passed with the same value as "node.id", but it is
deprecated since 4.4 and will not be passed
+ * from Apache Kafka 5.0 onwards. Implementations should read "node.id"
instead.
*
* Implement {@link org.apache.kafka.common.metrics.Monitorable} to enable the
manager to register metrics.
* The following tags are automatically added to all metrics registered:
<code>config</code> set to
diff --git
a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerConfig.java
b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerConfig.java
index ec9d9ba8263..d41d1a5a135 100644
---
a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerConfig.java
+++
b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerConfig.java
@@ -21,9 +21,13 @@ import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.config.ConfigDef;
+import org.apache.kafka.common.config.ConfigException;
import org.apache.kafka.common.serialization.ByteArrayDeserializer;
import org.apache.kafka.common.serialization.ByteArraySerializer;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
@@ -41,6 +45,8 @@ import static
org.apache.kafka.common.utils.internals.ConfigUtils.configMapToRed
*/
public final class TopicBasedRemoteLogMetadataManagerConfig {
+ private static final Logger LOG =
LoggerFactory.getLogger(TopicBasedRemoteLogMetadataManagerConfig.class);
+
public static final String REMOTE_LOG_METADATA_TOPIC_NAME =
"__remote_log_metadata";
public static final String
REMOTE_LOG_METADATA_TOPIC_REPLICATION_FACTOR_PROP =
"remote.log.metadata.topic.replication.factor";
@@ -81,7 +87,12 @@ public final class TopicBasedRemoteLogMetadataManagerConfig {
public static final String REMOTE_LOG_METADATA_PRODUCER_PREFIX =
"remote.log.metadata.producer.";
public static final String REMOTE_LOG_METADATA_CONSUMER_PREFIX =
"remote.log.metadata.consumer.";
public static final String REMOTE_LOG_METADATA_ADMIN_PREFIX =
"remote.log.metadata.admin.";
+ /**
+ * @deprecated Use {@link #NODE_ID} instead. This key is no longer passed
to plugins from Kafka 5.0 (KIP-1232).
+ */
+ @Deprecated(since = "4.4", forRemoval = true)
public static final String BROKER_ID = "broker.id";
+ public static final String NODE_ID = "node.id";
public static final String LOG_DIR = "log.dir";
private static final String REMOTE_LOG_METADATA_CLIENT_PREFIX =
"__remote_log_metadata_client";
@@ -138,10 +149,24 @@ public final class
TopicBasedRemoteLogMetadataManagerConfig {
consumeWaitMs = (long)
parsedConfigs.get(REMOTE_LOG_METADATA_CONSUME_WAIT_MS_PROP);
initializationRetryIntervalMs = (long)
parsedConfigs.get(REMOTE_LOG_METADATA_INITIALIZATION_RETRY_INTERVAL_MS_PROP);
initializationRetryMaxTimeoutMs = (long)
parsedConfigs.get(REMOTE_LOG_METADATA_INITIALIZATION_RETRY_MAX_TIMEOUT_MS_PROP);
- clientIdPrefix = REMOTE_LOG_METADATA_CLIENT_PREFIX + "_" +
props.get(BROKER_ID);
+ clientIdPrefix = REMOTE_LOG_METADATA_CLIENT_PREFIX + "_" +
nodeId(props);
initializeClientProperties(props);
}
+ private static Object nodeId(Map<String, ?> props) {
+ Object nodeId = props.get(NODE_ID);
+ if (nodeId != null) {
+ return nodeId;
+ }
+ Object brokerId = props.get(BROKER_ID);
+ if (brokerId == null) {
+ throw new ConfigException("Both " + NODE_ID + " and " + BROKER_ID
+ " configs are missing. Please configure " + NODE_ID + ".");
+ }
+ LOG.warn("The '{}' config is deprecated and will no longer be read in
Apache Kafka 5.0. Please use '{}' instead.",
+ BROKER_ID, NODE_ID);
+ return brokerId;
+ }
+
private void initializeClientProperties(Map<String, ?> configs) {
Map<String, Object> commonClientConfigs = new HashMap<>();
Map<String, Object> producerOnlyConfigs = new HashMap<>();
diff --git
a/storage/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogManager.java
b/storage/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogManager.java
index 6ec52ca48f2..1ac6028af51 100644
---
a/storage/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogManager.java
+++
b/storage/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogManager.java
@@ -47,7 +47,6 @@ import org.apache.kafka.common.utils.internals.ThreadUtils;
import org.apache.kafka.server.common.CheckpointFile;
import org.apache.kafka.server.common.OffsetAndEpoch;
import org.apache.kafka.server.common.StopPartition;
-import org.apache.kafka.server.config.ServerConfigs;
import org.apache.kafka.server.log.remote.TopicPartitionLog;
import
org.apache.kafka.server.log.remote.metadata.storage.ClassLoaderAwareRemoteLogMetadataManager;
import org.apache.kafka.server.log.remote.quota.RLMQuotaManager;
@@ -153,6 +152,10 @@ public class RemoteLogManager implements Closeable,
AsyncOffsetReader {
private static final Logger LOGGER =
LoggerFactory.getLogger(RemoteLogManager.class);
private static final String REMOTE_LOG_READER_THREAD_NAME_PATTERN =
"remote-log-reader-%d";
+ // The keys the RSM/RLMM plugins are configured with, independent of the
broker configs of the same
+ // names. broker.id is deprecated and will no longer be passed from 5.0
(KIP-1232).
+ private static final String BROKER_ID = "broker.id";
+ private static final String NODE_ID = "node.id";
private final RemoteLogManagerConfig rlmConfig;
private final int brokerId;
private final String logDir;
@@ -401,7 +404,8 @@ public class RemoteLogManager implements Closeable,
AsyncOffsetReader {
private Plugin<RemoteStorageManager>
configAndWrapRsmPlugin(RemoteStorageManager rsm) {
final Map<String, Object> rsmProps = new
HashMap<>(rlmConfig.remoteStorageManagerProps());
- rsmProps.put(ServerConfigs.BROKER_ID_CONFIG, brokerId);
+ rsmProps.put(BROKER_ID, brokerId);
+ rsmProps.put(NODE_ID, brokerId);
rsm.configure(rsmProps);
return Plugin.wrapInstance(rsm, metrics,
RemoteLogManagerConfig.REMOTE_STORAGE_MANAGER_CLASS_NAME_PROP);
}
@@ -428,7 +432,8 @@ public class RemoteLogManager implements Closeable,
AsyncOffsetReader {
// update the remoteLogMetadataProps here to override endpoint config
if any
rlmmProps.putAll(rlmConfig.remoteLogMetadataManagerProps());
- rlmmProps.put(ServerConfigs.BROKER_ID_CONFIG, brokerId);
+ rlmmProps.put(BROKER_ID, brokerId);
+ rlmmProps.put(NODE_ID, brokerId);
rlmmProps.put(LOG_DIR_CONFIG, logDir);
rlmmProps.put("cluster.id", clusterId);
diff --git
a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/RemoteLogMetadataManagerTestUtils.java
b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/RemoteLogMetadataManagerTestUtils.java
index fb8812d6b79..7108ccc0bca 100644
---
a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/RemoteLogMetadataManagerTestUtils.java
+++
b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/RemoteLogMetadataManagerTestUtils.java
@@ -25,8 +25,8 @@ import java.util.Objects;
import java.util.function.Function;
import java.util.function.Supplier;
-import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.BROKER_ID;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.LOG_DIR;
+import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.NODE_ID;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_COMMON_CLIENT_PREFIX;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_PARTITIONS_PROP;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_REPLICATION_FACTOR_PROP;
@@ -80,7 +80,7 @@ public class RemoteLogMetadataManagerTestUtils {
// Initialize TopicBasedRemoteLogMetadataManager.
Map<String, Object> configs = new HashMap<>();
configs.put(REMOTE_LOG_METADATA_COMMON_CLIENT_PREFIX +
CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
- configs.put(BROKER_ID, 0);
+ configs.put(NODE_ID, 0);
configs.put(LOG_DIR, logDir);
configs.put(REMOTE_LOG_METADATA_TOPIC_PARTITIONS_PROP,
METADATA_TOPIC_PARTITIONS_COUNT);
configs.put(REMOTE_LOG_METADATA_TOPIC_REPLICATION_FACTOR_PROP,
METADATA_TOPIC_REPLICATION_FACTOR);
diff --git
a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerConfigTest.java
b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerConfigTest.java
index bad522ff51e..e59d9f585e3 100644
---
a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerConfigTest.java
+++
b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerConfigTest.java
@@ -20,6 +20,7 @@ import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.common.config.ConfigException;
import org.apache.kafka.common.config.SaslConfigs;
import org.apache.kafka.common.config.SslConfigs;
import org.apache.kafka.test.TestUtils;
@@ -35,6 +36,7 @@ import java.util.Map;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.BROKER_ID;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.DEFAULT_REMOTE_LOG_METADATA_TOPIC_MIN_ISR;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.LOG_DIR;
+import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.NODE_ID;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_ADMIN_PREFIX;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_COMMON_CLIENT_PREFIX;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_CONSUMER_PREFIX;
@@ -44,6 +46,7 @@ import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemo
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_REPLICATION_FACTOR_PROP;
import static
org.apache.kafka.server.log.remote.metadata.storage.TopicBasedRemoteLogMetadataManagerConfig.REMOTE_LOG_METADATA_TOPIC_RETENTION_MS_PROP;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class TopicBasedRemoteLogMetadataManagerConfigTest {
@@ -174,7 +177,7 @@ public class TopicBasedRemoteLogMetadataManagerConfigTest {
Map<String, Object>
adminConfig) {
Map<String, Object> props = new HashMap<>();
props.put(REMOTE_LOG_METADATA_COMMON_CLIENT_PREFIX +
CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVERS);
- props.put(BROKER_ID, 1);
+ props.put(NODE_ID, 1);
props.put(LOG_DIR, TestUtils.tempDirectory().getAbsolutePath());
props.put(REMOTE_LOG_METADATA_TOPIC_REPLICATION_FACTOR_PROP, (short)
3);
props.put(REMOTE_LOG_METADATA_TOPIC_PARTITIONS_PROP, 10);
@@ -241,4 +244,30 @@ public class TopicBasedRemoteLogMetadataManagerConfigTest {
TopicBasedRemoteLogMetadataManagerConfig rlmmConfig = new
TopicBasedRemoteLogMetadataManagerConfig(props);
assertEquals(customMinIsr, rlmmConfig.metadataTopicMinIsr());
}
-}
\ No newline at end of file
+
+ @SuppressWarnings("removal")
+ @Test
+ public void testDeprecatedBrokerIdIsStillAccepted() {
+ Map<String, Object> props = createValidConfigProps();
+ props.remove(NODE_ID);
+ props.put(BROKER_ID, 1);
+ TopicBasedRemoteLogMetadataManagerConfig rlmmConfig = new
TopicBasedRemoteLogMetadataManagerConfig(props);
+
assertTrue(rlmmConfig.consumerProperties().get(CommonClientConfigs.CLIENT_ID_CONFIG).toString().endsWith("_1_consumer"));
+ }
+
+ @SuppressWarnings("removal")
+ @Test
+ public void testNodeIdTakesPrecedenceOverDeprecatedBrokerId() {
+ Map<String, Object> props = createValidConfigProps();
+ props.put(BROKER_ID, 2);
+ TopicBasedRemoteLogMetadataManagerConfig rlmmConfig = new
TopicBasedRemoteLogMetadataManagerConfig(props);
+
assertTrue(rlmmConfig.consumerProperties().get(CommonClientConfigs.CLIENT_ID_CONFIG).toString().endsWith("_1_consumer"));
+ }
+
+ @Test
+ public void testMissingNodeIdThrows() {
+ Map<String, Object> props = createValidConfigProps();
+ props.remove(NODE_ID);
+ assertThrows(ConfigException.class, () -> new
TopicBasedRemoteLogMetadataManagerConfig(props));
+ }
+}
diff --git
a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java
b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java
index 77233e15a2a..34c7b34473a 100644
---
a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java
+++
b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java
@@ -360,7 +360,7 @@ public class TopicBasedRemoteLogMetadataManagerTest {
// configure rlmm without bootstrap servers, so it will fail to
initialize admin client.
Map<String, Object> configs = Map.of(
TopicBasedRemoteLogMetadataManagerConfig.LOG_DIR,
TestUtils.tempDirectory("rlmm_segs_").getAbsolutePath(),
- TopicBasedRemoteLogMetadataManagerConfig.BROKER_ID, 0
+ TopicBasedRemoteLogMetadataManagerConfig.NODE_ID, 0
);
rlmm.configure(configs);
rlmm.onBrokerReady();
diff --git
a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/LocalTieredStorageTest.java
b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/LocalTieredStorageTest.java
index a11ba3a2423..bffa2fb2f37 100644
---
a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/LocalTieredStorageTest.java
+++
b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/LocalTieredStorageTest.java
@@ -91,7 +91,7 @@ public final class LocalTieredStorageTest {
Map<String, Object> config = new HashMap<>();
config.put(LocalTieredStorage.STORAGE_DIR_CONFIG, storageDir);
config.put(LocalTieredStorage.DELETE_ON_CLOSE_CONFIG, "true");
- config.put(LocalTieredStorage.BROKER_ID, 1);
+ config.put(LocalTieredStorage.NODE_ID, 1);
config.putAll(extraConfig);
tieredStorage.configure(config);
diff --git
a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/RemoteLogManagerTest.java
b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/RemoteLogManagerTest.java
index b6d5ae7aefb..9f1f15a8872 100644
---
a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/RemoteLogManagerTest.java
+++
b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/RemoteLogManagerTest.java
@@ -39,7 +39,6 @@ import org.apache.kafka.common.utils.MockTime;
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.server.common.OffsetAndEpoch;
import org.apache.kafka.server.common.StopPartition;
-import org.apache.kafka.server.config.ServerConfigs;
import org.apache.kafka.server.log.remote.TopicPartitionLog;
import org.apache.kafka.server.log.remote.quota.RLMQuotaManager;
import org.apache.kafka.server.log.remote.quota.RLMQuotaManagerConfig;
@@ -385,7 +384,8 @@ public class RemoteLogManagerTest {
assertEquals(host + ":" + port,
capture.getValue().get(REMOTE_LOG_METADATA_COMMON_CLIENT_PREFIX +
"bootstrap.servers"));
assertEquals(securityProtocol,
capture.getValue().get(REMOTE_LOG_METADATA_COMMON_CLIENT_PREFIX +
"security.protocol"));
assertEquals(clusterId, capture.getValue().get("cluster.id"));
- assertEquals(brokerId,
capture.getValue().get(ServerConfigs.BROKER_ID_CONFIG));
+ assertEquals(brokerId, capture.getValue().get("broker.id"));
+ assertEquals(brokerId, capture.getValue().get("node.id"));
}
@SuppressWarnings("unchecked")
@@ -423,7 +423,8 @@ public class RemoteLogManagerTest {
// should be overridden as SSL
assertEquals("SSL",
capture.getValue().get(REMOTE_LOG_METADATA_COMMON_CLIENT_PREFIX +
"security.protocol"));
assertEquals(clusterId, capture.getValue().get("cluster.id"));
- assertEquals(brokerId,
capture.getValue().get(ServerConfigs.BROKER_ID_CONFIG));
+ assertEquals(brokerId, capture.getValue().get("broker.id"));
+ assertEquals(brokerId, capture.getValue().get("node.id"));
}
}
@@ -433,10 +434,12 @@ public class RemoteLogManagerTest {
ArgumentCaptor<Map<String, Object>> capture =
ArgumentCaptor.forClass(Map.class);
verify(remoteStorageManager, times(1)).configure(capture.capture());
assertEquals(brokerId, capture.getValue().get("broker.id"));
+ assertEquals(brokerId, capture.getValue().get("node.id"));
assertEquals(remoteLogStorageTestVal,
capture.getValue().get(remoteLogStorageTestProp));
verify(remoteLogMetadataManager,
times(1)).configure(capture.capture());
assertEquals(brokerId, capture.getValue().get("broker.id"));
+ assertEquals(brokerId, capture.getValue().get("node.id"));
assertEquals(logDir, capture.getValue().get("log.dir"));
// verify the configs starting with "remote.log.metadata",
"remote.log.metadata.common.client."
diff --git
a/storage/src/testFixtures/java/org/apache/kafka/server/log/remote/storage/LocalTieredStorage.java
b/storage/src/testFixtures/java/org/apache/kafka/server/log/remote/storage/LocalTieredStorage.java
index 926a532826e..76e925c3639 100644
---
a/storage/src/testFixtures/java/org/apache/kafka/server/log/remote/storage/LocalTieredStorage.java
+++
b/storage/src/testFixtures/java/org/apache/kafka/server/log/remote/storage/LocalTieredStorage.java
@@ -132,8 +132,14 @@ public final class LocalTieredStorage implements
RemoteStorageManager {
public static final String ENABLE_DELETE_API_CONFIG = "delete.enable";
/**
- * The ID of the broker which owns this instance of {@link
LocalTieredStorage}.
+ * The ID of the node which owns this instance of {@link
LocalTieredStorage}.
*/
+ public static final String NODE_ID = "node.id";
+
+ /**
+ * @deprecated Use {@link #NODE_ID} instead. This key is no longer read
from Kafka 5.0 (KIP-1232).
+ */
+ @Deprecated(since = "4.4", forRemoval = true)
public static final String BROKER_ID = "broker.id";
private static final String ROOT_STORAGE_DIR_NAME = "kafka-tiered-storage";
@@ -221,6 +227,19 @@ public final class LocalTieredStorage implements
RemoteStorageManager {
this.storageListeners.add(listener);
}
+ private Integer nodeId(final Map<String, ?> configs) {
+ final Integer nodeId = (Integer) configs.get(NODE_ID);
+ if (nodeId != null) {
+ return nodeId;
+ }
+ final Integer brokerId = (Integer) configs.get(BROKER_ID);
+ if (brokerId != null) {
+ logger.warn("The '{}' config is deprecated and will no longer be
read in Apache Kafka 5.0. Please use '{}' instead.",
+ BROKER_ID, NODE_ID);
+ }
+ return brokerId;
+ }
+
@Override
public void configure(Map<String, ?> configs) {
if (storageDirectory != null) {
@@ -233,14 +252,15 @@ public final class LocalTieredStorage implements
RemoteStorageManager {
final String shouldDeleteOnClose = (String)
configs.get(DELETE_ON_CLOSE_CONFIG);
final String transfererClass = (String)
configs.get(TRANSFERER_CLASS_CONFIG);
final String isDeleteEnabled = (String)
configs.get(ENABLE_DELETE_API_CONFIG);
- final Integer brokerIdInt = (Integer) configs.get(BROKER_ID);
+ final Integer nodeIdInt = nodeId(configs);
- if (brokerIdInt == null) {
- throw new InvalidConfigurationException(
- "Broker ID is required to configure the LocalTieredStorage
manager.");
+ if (nodeIdInt == null) {
+ throw new InvalidConfigurationException(format(
+ "Both %s and %s configs are missing. Please configure %s
to use the LocalTieredStorage manager.",
+ NODE_ID, BROKER_ID, NODE_ID));
}
- brokerId = brokerIdInt;
+ brokerId = nodeIdInt;
logger = new LogContext(format("[LocalTieredStorage Id=%d] ",
brokerId)).logger(this.getClass());
if (shouldDeleteOnClose != null) {