From 174ef3846a4c65dd10b58e2e92f9f625cb208950 Mon Sep 17 00:00:00 2001 From: chickenchickenlove Date: Thu, 6 Aug 2026 22:48:44 +0900 Subject: [PATCH] KAFKA-20889: reject empty group.instance.id values in Kafka Streams. --- .../org/apache/kafka/streams/StreamsConfig.java | 10 +++++++++- .../org/apache/kafka/streams/StreamsConfigTest.java | 13 +++++++++++++ 2 files changed, 22 insertions(+), 1 deletion(-) 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 bc33f7308690d..55ee36be36c80 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -1933,7 +1933,15 @@ public Map getMainConsumerConfigs(final String groupId, final St // add group id, client id with stream client id prefix, and group instance id consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); consumerProps.put(CommonClientConfigs.CLIENT_ID_CONFIG, clientId); - final String groupInstanceId = (String) consumerProps.get(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG); + final String groupInstanceId = (String) parseType( + ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, + consumerProps.get(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG), + Type.STRING + ); + new ConfigDef.NonEmptyString().ensureValid( + ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, + groupInstanceId + ); // Suffix each thread consumer with thread.id to enforce uniqueness of group.instance.id. if (groupInstanceId != null) { consumerProps.put(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, groupInstanceId + "-" + threadIdx); 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 87648a25171ad..35cdd168e4cd2 100644 --- a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java @@ -243,6 +243,19 @@ public void testGetGroupInstanceIdConfigs() { assertNull(returnedProps.get(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG)); } + @ParameterizedTest + @ValueSource(strings = {"", StreamsConfig.CONSUMER_PREFIX, StreamsConfig.MAIN_CONSUMER_PREFIX}) + public void shouldRejectEmptyGroupInstanceId(final String prefix) { + props.put(prefix + ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, ""); + final StreamsConfig streamsConfig = new StreamsConfig(props); + + final ConfigException exception = assertThrows( + ConfigException.class, + () -> streamsConfig.getMainConsumerConfigs(groupId, clientId, threadIdx) + ); + assertTrue(exception.getMessage().contains(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG)); + } + @ParameterizedTest @ValueSource(strings = {"", StreamsConfig.CONSUMER_PREFIX, StreamsConfig.MAIN_CONSUMER_PREFIX}) public void shouldAllowStaticMembershipWhenStreamsProtocolUsed(final String prefix) {