From 6385fd45dcc19869e440c71cc451523348184887 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=EC=97=84=EC=9C=A4=EC=84=AD?= <62834176+lh0156@users.noreply.github.com> Date: Sun, 2 Aug 2026 22:39:11 +0900 Subject: [PATCH] KAFKA-20869: Guard stream thread addition for global-only topologies Prevent addStreamThread from creating a local stream thread when the topology has no local processing tasks. Add a regression test covering a running global-only client. Generated-by: OpenAI Codex (GPT-5) --- .../org/apache/kafka/streams/KafkaStreams.java | 5 +++++ .../apache/kafka/streams/KafkaStreamsTest.java | 16 ++++++++++++++++ 2 files changed, 21 insertions(+) diff --git a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java index aecdf3f67bbfe..974fc473586ee 100644 --- a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java +++ b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java @@ -1136,6 +1136,11 @@ private static Metrics createMetrics(final StreamsConfig config, final Time time * @return name of the added stream thread or empty if a new stream thread could not be added */ public Optional addStreamThread() { + if (topologyMetadata.hasNoLocalTopology()) { + log.warn("Cannot add a stream thread to a topology with no local processing tasks"); + return Optional.empty(); + } + if (isRunningOrRebalancing()) { final StreamThread streamThread; synchronized (changeThreadCount) { diff --git a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java index 528efd03dbc8c..5bb1c52c2ec64 100644 --- a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java @@ -721,6 +721,22 @@ public void shouldAddThreadWhenRunning() throws Exception { } } + @Test + public void shouldNotAddThreadToGlobalOnlyTopology() { + prepareStreams(); + final StreamsBuilder builder = new StreamsBuilder(); + builder.globalTable("global-topic"); + props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 1); + + try (final KafkaStreams streams = new KafkaStreams(builder.build(), props, supplier, time)) { + streams.start(); + + assertThat(streams.threads.size(), equalTo(0)); + assertThat(streams.addStreamThread(), equalTo(Optional.empty())); + assertThat(streams.threads.size(), equalTo(0)); + } + } + @Test public void shouldNotAddThreadWhenCreated() { prepareStreams();