This is an automated email from the ASF dual-hosted git repository.
bbejeck 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 f324cb59591 MINOR: Set StreamThread to non-daemon status (#22902)
f324cb59591 is described below
commit f324cb595910d42a402daaa3fa18f44bd0534765
Author: Bill Bejeck <[email protected]>
AuthorDate: Thu Jul 23 14:34:02 2026 -0400
MINOR: Set StreamThread to non-daemon status (#22902)
This PR sets
1. `StreamThread` `daemon=false`
2. `GlobalStreamThread` `daemon=false`
3. `StateUpdaterThread` `daemon=true`
Reviewers: Matthias J. Sax <[email protected]>
---
streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java | 4 ++++
.../apache/kafka/streams/processor/internals/DefaultStateUpdater.java | 2 ++
.../apache/kafka/streams/processor/internals/GlobalStreamThread.java | 3 +++
.../org/apache/kafka/streams/processor/internals/StreamThread.java | 2 ++
4 files changed, 11 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 0c2b5582fa6..ad9aada7c0d 100644
--- a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java
+++ b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java
@@ -1370,6 +1370,7 @@ public class KafkaStreams implements AutoCloseable {
private ScheduledExecutorService setupStateDirCleaner() {
return Executors.newSingleThreadScheduledExecutor(r -> {
final Thread thread = new Thread(r, clientId + "-CleanupThread");
+ // daemon: background cleanup must not block JVM shutdown
thread.setDaemon(true);
return thread;
});
@@ -1380,6 +1381,7 @@ public class KafkaStreams implements AutoCloseable {
if
(RecordingLevel.forName(config.getString(METRICS_RECORDING_LEVEL_CONFIG)) ==
RecordingLevel.DEBUG) {
return Executors.newSingleThreadScheduledExecutor(r -> {
final Thread thread = new Thread(r, clientId +
"-RocksDBMetricsRecordingTrigger");
+ // daemon: background metrics recording must not block JVM
shutdown
thread.setDaemon(true);
return thread;
});
@@ -1620,6 +1622,7 @@ public class KafkaStreams implements AutoCloseable {
final Thread shutdownThread = shutdownHelper(false, timeoutMs,
operation);
+ // daemon: the shutdown thread itself must not prevent JVM exit
shutdownThread.setDaemon(true);
shutdownThread.start();
@@ -1638,6 +1641,7 @@ public class KafkaStreams implements AutoCloseable {
} else {
final Thread shutdownThread = shutdownHelper(true, -1,
org.apache.kafka.streams.CloseOptions.GroupMembershipOperation.REMAIN_IN_GROUP);
+ // daemon: the shutdown thread itself must not prevent JVM exit
shutdownThread.setDaemon(true);
shutdownThread.start();
}
diff --git
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java
index bdde0c44435..68e86449bda 100644
---
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java
+++
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java
@@ -109,6 +109,8 @@ public class DefaultStateUpdater implements StateUpdater {
final StreamsMetricsImpl metrics,
final ChangelogReader changelogReader) {
super(name);
+ // daemon: internal helper thread must not block JVM shutdown on
its own
+ setDaemon(true);
this.changelogReader = changelogReader;
this.updaterMetrics = new StateUpdaterMetrics(metrics, name);
this.metricsConfig = metrics.metricsRegistry().config();
diff --git
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java
index 4b5bc646e0d..2e93941f527 100644
---
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java
+++
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java
@@ -198,6 +198,7 @@ public class GlobalStreamThread extends Thread {
}
}
+ @SuppressWarnings("this-escape")
public GlobalStreamThread(final ProcessorTopology topology,
final StreamsConfig config,
final Consumer<byte[], byte[]> globalConsumer,
@@ -210,6 +211,8 @@ public class GlobalStreamThread extends Thread {
final StateRestoreListener stateRestoreListener,
final java.util.function.Consumer<Throwable>
streamsUncaughtExceptionHandler) {
super(threadClientId);
+ // explicitly non-daemon so the JVM doesn't exit while this thread is
still restoring/serving global state
+ setDaemon(false);
this.time = time;
this.config = config;
this.topology = topology;
diff --git
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java
index d731215f9d3..dcd853e2a88 100644
---
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java
+++
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java
@@ -828,6 +828,8 @@ public class StreamThread extends Thread implements
ProcessingThread {
final long maxUncommittedBytesPerThread
) {
super(threadId);
+ // explicitly non-daemon so the JVM doesn't exit while this thread is
still processing
+ setDaemon(false);
this.stateLock = new Object();
this.adminClient = adminClient;
this.streamsMetrics = streamsMetrics;