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