This is an automated email from the ASF dual-hosted git repository.

mjsax 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 1b5d51612b2 KAFKA-20403 : streams - Fix stream threads interruptions 
(#21970)
1b5d51612b2 is described below

commit 1b5d51612b29ad2b11c6711dc1fe447db1c2d938
Author: Murali Basani <[email protected]>
AuthorDate: Thu Jul 16 01:23:00 2026 +0200

    KAFKA-20403 : streams - Fix stream threads interruptions (#21970)
    
    This PR resets thread interrupted signals and improves corresponding
    logging.
    
    Reviewers: Bill Bejeck <[email protected]>, Apoorv Mittal
    <[email protected]>, Matthias J. Sax <[email protected]>
---
 .../kafka/streams/processor/internals/DefaultStateUpdater.java | 10 +++++++---
 .../kafka/streams/processor/internals/TopologyMetadata.java    |  4 +++-
 .../internals/namedtopology/AddNamedTopologyResult.java        |  1 +
 .../namedtopology/KafkaStreamsNamedTopologyWrapper.java        |  6 ++++--
 .../streams/processor/internals/tasks/DefaultTaskExecutor.java | 10 +++++++---
 .../streams/processor/internals/tasks/DefaultTaskManager.java  |  7 +++----
 6 files changed, 25 insertions(+), 13 deletions(-)

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 dfd8463fc2f..bdde0c44435 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
@@ -492,8 +492,8 @@ public class DefaultStateUpdater implements StateUpdater {
                     tasksAndActionsCondition.await();
                 }
             } catch (final InterruptedException ignored) {
-                // we never interrupt the thread, but only signal the condition
-                // and hence this exception should never be thrown
+                Thread.currentThread().interrupt();
+                log.warn("State updater thread was interrupted while waiting 
for changelogs");
             } finally {
                 tasksAndActionsLock.unlock();
                 isIdle.set(false);
@@ -922,6 +922,8 @@ public class DefaultStateUpdater implements StateUpdater {
                 }
                 stateUpdaterThread = null;
             } catch (final InterruptedException ignored) {
+                Thread.currentThread().interrupt();
+                log.warn("Interrupted while waiting for state updater thread 
to shut down");
             }
         }
     }
@@ -999,7 +1001,9 @@ public class DefaultStateUpdater implements StateUpdater {
                 now = time.milliseconds();
             }
             return result;
-        } catch (final InterruptedException ignored) {
+        } catch (final InterruptedException e) {
+            Thread.currentThread().interrupt();
+            log.warn("StateUpdaterThread was interrupted while waiting", e);
         }
         return result;
     }
diff --git 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java
 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java
index 6b66a08333d..353a7135beb 100644
--- 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java
+++ 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TopologyMetadata.java
@@ -230,7 +230,9 @@ public class TopologyMetadata {
                         log.debug("Detected that the topology is currently 
empty, waiting for something to process");
                         version.topologyCV.await();
                     } catch (final InterruptedException e) {
-                        log.error("StreamThread was interrupted while waiting 
on empty topology", e);
+                        Thread.currentThread().interrupt();
+                        log.warn("StreamThread was interrupted while waiting 
on empty topology", e);
+                        break;
                     }
                 }
             } finally {
diff --git 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/AddNamedTopologyResult.java
 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/AddNamedTopologyResult.java
index 048eacda5b6..17c2fddc0ba 100644
--- 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/AddNamedTopologyResult.java
+++ 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/AddNamedTopologyResult.java
@@ -54,6 +54,7 @@ public class AddNamedTopologyResult {
                 return new StreamsException(e.getCause());
             }
         } catch (final InterruptedException e) {
+            Thread.currentThread().interrupt();
             return null;
         }
     }
diff --git 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/KafkaStreamsNamedTopologyWrapper.java
 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/KafkaStreamsNamedTopologyWrapper.java
index 655a7d46565..b0a2a3ff07c 100644
--- 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/KafkaStreamsNamedTopologyWrapper.java
+++ 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/namedtopology/KafkaStreamsNamedTopologyWrapper.java
@@ -306,7 +306,7 @@ public class KafkaStreamsNamedTopologyWrapper extends 
KafkaStreams {
                 log.info("Successfully completed resetting offsets.");
                 break;
             } catch (final InterruptedException ex) {
-                ex.printStackTrace();
+                Thread.currentThread().interrupt();
                 log.error("Offset reset failed.", ex);
                 throw new StreamsException(ex);
             } catch (final ExecutionException ex) {
@@ -331,7 +331,9 @@ public class KafkaStreamsNamedTopologyWrapper extends 
KafkaStreams {
             try {
                 Thread.sleep(100);
             } catch (final InterruptedException ex) {
-                ex.printStackTrace();
+                Thread.currentThread().interrupt();
+                log.warn("Interrupted during offset reset retry backoff", ex);
+                break;
             }
         }
     }
diff --git 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutor.java
 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutor.java
index ff127e2071b..fbb5f77dddd 100644
--- 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutor.java
+++ 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskExecutor.java
@@ -117,8 +117,10 @@ public class DefaultTaskExecutor implements TaskExecutor {
             if (currentTask == null) {
                 try {
                     taskManager.awaitProcessableTasks(shutdownRequested::get);
-                } catch (final InterruptedException ignored) {
-                    // Can be ignored, the cause of the interrupted will be 
handled in the event loop
+                } catch (final InterruptedException e) {
+                    Thread.currentThread().interrupt();
+                    // The event loop will check shutdownRequested on the next 
iteration
+                    log.warn("TaskExecutorThread was interrupted while 
waiting", e);
                 }
             } else {
                 boolean progressed = false;
@@ -266,7 +268,9 @@ public class DefaultTaskExecutor implements TaskExecutor {
                     throw new StreamsException("State updater thread did not 
shutdown within the timeout");
                 }
                 taskExecutorThread = null;
-            } catch (final InterruptedException ignored) {
+            } catch (final InterruptedException e) {
+                Thread.currentThread().interrupt();
+                log.warn("Interrupted while waiting for task executor thread 
to shut down", e);
             }
         }
     }
diff --git 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java
 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java
index 81a1bb02ff4..c4fab3b5d5d 100644
--- 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java
+++ 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/tasks/DefaultTaskManager.java
@@ -146,10 +146,9 @@ public final class DefaultTaskManager implements 
TaskManager {
                 } else {
                     log.debug("Not awaiting since shutdown was requested");
                 }
-            } catch (final InterruptedException ignored) {
-                // we interrupt the thread for shut down and pause.
-                // we can ignore this exception.
-                log.debug("Await unblocked: Interrupted while waiting for 
processable tasks");
+            } catch (final InterruptedException e) {
+                Thread.currentThread().interrupt();
+                log.warn("Await unblocked: Interrupted while waiting for 
processable tasks", e);
                 return true;
             }
             log.debug("Await unblocked: Woken up to check for processable 
tasks");

Reply via email to