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 d04a53943c691..84ed4fa6a86ba 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 @@ -719,7 +719,7 @@ void runOnce() { // Should only proceed when the thread is still running after #pollRequests(), because no external state mutation // could affect the task manager state beyond this point within #runOnce(). if (!isRunning()) { - log.debug("Thread state is already {}, skipping the run once call after poll request", state); + log.info("Thread state is already {}, skipping the run once call after poll request", state); return; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsRebalanceListener.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsRebalanceListener.java index ab8a5ff86e976..ba2883b3864ce 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsRebalanceListener.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsRebalanceListener.java @@ -87,7 +87,9 @@ public void onPartitionsRevoked(final Collection partitions) { taskManager.activeTaskIds(), taskManager.standbyTaskIds()); - if (streamThread.setState(State.PARTITIONS_REVOKED) != null && !partitions.isEmpty()) { + // We need to still invoke handleRevocation if the thread has been told to shut down, but we shouldn't ever + // transition away from PENDING_SHUTDOWN once it's been initiated (to anything other than DEAD) + if ((streamThread.setState(State.PARTITIONS_REVOKED) != null || streamThread.state() == State.PENDING_SHUTDOWN) && !partitions.isEmpty()) { final long start = time.milliseconds(); try { taskManager.handleRevocation(partitions);