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;

Reply via email to