This is an automated email from the ASF dual-hosted git repository.
mjsax 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 076ea535fa3 KAFKA-17164: Enforce 'application.server' <server>:<port>
format at config level. (#22202)
076ea535fa3 is described below
commit 076ea535fa39050d27c5ab7d8ea75361d1c4b0a2
Author: ChickenchickenLove <[email protected]>
AuthorDate: Wed Jun 24 02:21:54 2026 +0900
KAFKA-17164: Enforce 'application.server' <server>:<port> format at config
level. (#22202)
Implements KIP-1245.
Enforces the `application.server` host:port format during StreamsConfig
construction.
Reviewers: Matthias J. Sax <[email protected]>, Nikita Shupletsov
<[email protected]>, Nilesh Kumar <[email protected]>
---
docs/streams/upgrade-guide.md | 2 +
.../org/apache/kafka/streams/StreamsConfig.java | 2 +
.../ApplicationServerConfigValidator.java | 47 ++++++++++++++++++++++
.../org/apache/kafka/streams/state/HostInfo.java | 15 ++++++-
.../apache/kafka/streams/StreamsConfigTest.java | 20 ++++++++-
.../apache/kafka/streams/state/HostInfoTest.java | 8 ++++
6 files changed, 91 insertions(+), 3 deletions(-)
diff --git a/docs/streams/upgrade-guide.md b/docs/streams/upgrade-guide.md
index 387bef80920..0951515f865 100644
--- a/docs/streams/upgrade-guide.md
+++ b/docs/streams/upgrade-guide.md
@@ -69,6 +69,8 @@ Since 2.6.0 release, Kafka Streams depends on a RocksDB
version that requires Ma
Kafka Streams no longer emits a WARN from `KafkaStreams#cleanUp()` when the
application state directory cannot be deleted only because expected metadata
files remain, such as `kafka-streams-process-metadata` and/or `.lock`. In this
case, the local state cleanup is considered successful and the application
state directory may be retained. Users who require a full local reset including
persisted process metadata should manually delete the application state
directory after the Kafka Streams [...]
+Kafka Streams now validates the `application.server` configuration when
`StreamsConfig` is created. The value must be empty or a valid endpoint from
which Kafka Streams can parse both host and port, such as `host:port` or
`protocol://host:port`. Invalid values that may previously have failed later
during startup or assignment now fail earlier with a `ConfigException`. More
details can be found in
[KIP-1245](https://cwiki.apache.org/confluence/display/KAFKA/KIP-1245%3A+Enforce+%27applicat
[...]
+
## Streams API changes in 4.3.0
**Note:** Kafka Streams 4.3.0 contains a critical native memory leak in the
RocksDB state store layer
([KAFKA-20616](https://issues.apache.org/jira/browse/KAFKA-20616)). The
`ColumnFamilyOptions` for the offsets column family is not closed, and column
family handles can leak on close-path exceptions, which under cascading task
closes (e.g., rebalances or error-triggered recoveries) leads to unbounded
off-heap memory growth and eventual OOM. Users running Kafka Streams should
consider upg [...]
diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
index d9e1569fe36..d6099f4cec8 100644
--- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
+++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
@@ -44,6 +44,7 @@ import
org.apache.kafka.streams.errors.LogAndFailProcessingExceptionHandler;
import org.apache.kafka.streams.errors.ProcessingExceptionHandler;
import org.apache.kafka.streams.errors.ProductionExceptionHandler;
import org.apache.kafka.streams.errors.StreamsException;
+import org.apache.kafka.streams.internals.ApplicationServerConfigValidator;
import org.apache.kafka.streams.internals.StreamsConfigUtils;
import org.apache.kafka.streams.internals.UpgradeFromValues;
import org.apache.kafka.streams.kstream.SessionWindowedDeserializer;
@@ -1114,6 +1115,7 @@ public class StreamsConfig extends AbstractConfig {
.define(APPLICATION_SERVER_CONFIG,
Type.STRING,
"",
+ new ApplicationServerConfigValidator(),
Importance.LOW,
APPLICATION_SERVER_DOC)
.define(BUFFERED_RECORDS_PER_PARTITION_CONFIG,
diff --git
a/streams/src/main/java/org/apache/kafka/streams/internals/ApplicationServerConfigValidator.java
b/streams/src/main/java/org/apache/kafka/streams/internals/ApplicationServerConfigValidator.java
new file mode 100644
index 00000000000..55e4af8fdbe
--- /dev/null
+++
b/streams/src/main/java/org/apache/kafka/streams/internals/ApplicationServerConfigValidator.java
@@ -0,0 +1,47 @@
+/*
+ * 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.streams.internals;
+
+import org.apache.kafka.common.config.ConfigDef;
+import org.apache.kafka.common.config.ConfigException;
+import org.apache.kafka.common.utils.Utils;
+import org.apache.kafka.streams.state.HostInfo;
+
+public class ApplicationServerConfigValidator implements ConfigDef.Validator {
+
+ @Override
+ public void ensureValid(final String name, final Object value) {
+ if (!(value instanceof String)) {
+ throw new ConfigException(name + " must be a string");
+ }
+
+ final String endPoint = (String) value;
+ if (Utils.isBlank(endPoint)) {
+ return;
+ }
+ try {
+ HostInfo.buildFromEndpoint(endPoint);
+ } catch (final ConfigException e) {
+ throw new ConfigException(name, value, e.getMessage());
+ }
+ }
+
+ @Override
+ public String toString() {
+ return "A host:port pair, protocol://host:port, or an empty string";
+ }
+}
diff --git a/streams/src/main/java/org/apache/kafka/streams/state/HostInfo.java
b/streams/src/main/java/org/apache/kafka/streams/state/HostInfo.java
index e2e17df0e58..1393dc8667e 100644
--- a/streams/src/main/java/org/apache/kafka/streams/state/HostInfo.java
+++ b/streams/src/main/java/org/apache/kafka/streams/state/HostInfo.java
@@ -57,7 +57,20 @@ public class HostInfo {
}
final String host = getHost(endPoint);
- final Integer port = getPort(endPoint);
+ if (Utils.isBlank(host)) {
+ throw new ConfigException(
+ String.format("Error parsing host address %s. Expected format
host:port.", endPoint)
+ );
+ }
+
+ final Integer port;
+ try {
+ port = getPort(endPoint);
+ } catch (final NumberFormatException e) {
+ throw new ConfigException(
+ String.format("Error parsing host address %s. Expected format
host:port.", endPoint)
+ );
+ }
if (host == null || port == null) {
throw new ConfigException(
diff --git
a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java
b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java
index c8f358d9fe6..481c1ad6bae 100644
--- a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java
+++ b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java
@@ -251,7 +251,7 @@ public class StreamsConfigTest {
props.put(StreamsConfig.MAX_WARMUP_REPLICAS_CONFIG, 9);
props.put(StreamsConfig.PROBING_REBALANCE_INTERVAL_MS_CONFIG, 99_999L);
props.put(StreamsConfig.WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG,
7L);
- props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, "dummy:host");
+ props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, "dummy:8080");
props.put(StreamsConfig.topicPrefix(TopicConfig.SEGMENT_BYTES_CONFIG),
1024 * 1024);
final StreamsConfig streamsConfig = new StreamsConfig(props);
final Map<String, Object> returnedProps =
streamsConfig.getMainConsumerConfigs(groupId, clientId, threadIdx);
@@ -266,7 +266,7 @@ public class StreamsConfigTest {
returnedProps.get(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG)
);
assertEquals(7L,
returnedProps.get(StreamsConfig.WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG));
- assertEquals("dummy:host",
returnedProps.get(StreamsConfig.APPLICATION_SERVER_CONFIG));
+ assertEquals("dummy:8080",
returnedProps.get(StreamsConfig.APPLICATION_SERVER_CONFIG));
assertEquals(1024 * 1024,
returnedProps.get(StreamsConfig.topicPrefix(TopicConfig.SEGMENT_BYTES_CONFIG)));
}
@@ -1902,6 +1902,22 @@ public class StreamsConfigTest {
}
}
+ @ParameterizedTest
+ @ValueSource(strings = {"dummy:host", "dummy:9999999999999999999999999",
"dummy", "dummy:", ":port", ":", ":8080"})
+ public void
shouldThrowConfigExceptionWithInvalidApplicationServerConfigValue(final String
applicationServerConfig) {
+ props.put(StreamsConfig.APPLICATION_SERVER_CONFIG,
applicationServerConfig);
+ assertThrows(ConfigException.class, () -> new StreamsConfig(props));
+ }
+
+ @ParameterizedTest
+ @ValueSource(strings = {"", "127.0.0.1:8080", "localhost:8080",
"[::1]:8080", "http://localhost:8080"})
+ public void shouldAcceptWithValidApplicationServerConfigValue(final String
applicationServerConfigValue) {
+ props.put(StreamsConfig.APPLICATION_SERVER_CONFIG,
applicationServerConfigValue);
+ final StreamsConfig streamsConfig = new StreamsConfig(props);
+ final Map<String, Object> returnedProps =
streamsConfig.getMainConsumerConfigs(groupId, clientId, threadIdx);
+ assertEquals(applicationServerConfigValue,
returnedProps.get(StreamsConfig.APPLICATION_SERVER_CONFIG));
+ }
+
@SuppressWarnings("deprecation")
@Test
public void
shouldNotLogWarningWhenProcessingExceptionHandlerIsEnabledOnGlobalThread() {
diff --git
a/streams/src/test/java/org/apache/kafka/streams/state/HostInfoTest.java
b/streams/src/test/java/org/apache/kafka/streams/state/HostInfoTest.java
index 10f7749f6a6..26dfeac7c24 100644
--- a/streams/src/test/java/org/apache/kafka/streams/state/HostInfoTest.java
+++ b/streams/src/test/java/org/apache/kafka/streams/state/HostInfoTest.java
@@ -19,6 +19,8 @@ package org.apache.kafka.streams.state;
import org.apache.kafka.common.config.ConfigException;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.is;
@@ -50,4 +52,10 @@ public class HostInfoTest {
public void shouldThrowConfigExceptionForNonsenseEndPoint() {
assertThrows(ConfigException.class, () ->
HostInfo.buildFromEndpoint("nonsense"));
}
+
+ @ParameterizedTest
+ @ValueSource(strings = {"dummy:host", "dummy:9999999999999999999999999",
"dummy", "dummy:", ":port", ":", ":8080"})
+ public void shouldThrowConfigExceptionWithInvalidEndpoint(final String
invalidEndpoint) {
+ assertThrows(ConfigException.class, () ->
HostInfo.buildFromEndpoint(invalidEndpoint));
+ }
}