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