From 572109605bf7548ac0474b3596b71b74a1c15395 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Fri, 12 Mar 2021 15:57:36 -0800 Subject: [PATCH 1/2] fix in StreamsRebalanceListener, exit early in StreamThread main loop after poll phase --- .../kafka/streams/processor/internals/StreamThread.java | 7 +++++++ .../processor/internals/StreamsRebalanceListener.java | 4 +++- 2 files changed, 10 insertions(+), 1 deletion(-) 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..821475e8e05f7 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 @@ -714,6 +714,13 @@ void runOnce() { final long pollLatency = pollPhase(); + // Optimization to skip the rest of the processing loop in case the thread was requested to shut down during + // the poll phase + if (!isRunning()) { + log.info("Exiting the processing loop early since StreamThread state is {}", state); + return; + } + // Shutdown hook could potentially be triggered and transit the thread state to PENDING_SHUTDOWN during #pollRequests(). // The task manager internal states could be uninitialized if the state transition happens during #onPartitionsAssigned(). // Should only proceed when the thread is still running after #pollRequests(), because no external state mutation 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); From 513ef98893a8c8e8216e410a53ec694d2ad9a2f1 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Fri, 12 Mar 2021 16:14:36 -0800 Subject: [PATCH 2/2] remove duplicated exit --- .../kafka/streams/processor/internals/StreamThread.java | 9 +-------- 1 file changed, 1 insertion(+), 8 deletions(-) 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 821475e8e05f7..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 @@ -714,19 +714,12 @@ void runOnce() { final long pollLatency = pollPhase(); - // Optimization to skip the rest of the processing loop in case the thread was requested to shut down during - // the poll phase - if (!isRunning()) { - log.info("Exiting the processing loop early since StreamThread state is {}", state); - return; - } - // Shutdown hook could potentially be triggered and transit the thread state to PENDING_SHUTDOWN during #pollRequests(). // The task manager internal states could be uninitialized if the state transition happens during #onPartitionsAssigned(). // 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; }