From d20095edbbe2f5f73bb250e094408f04a7923039 Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Fri, 1 Feb 2019 11:16:37 -0800 Subject: [PATCH 1/9] If the thread is not alive, complete shutdown when calling `shutdown()`. --- core/src/main/scala/kafka/utils/ShutdownableThread.scala | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/utils/ShutdownableThread.scala b/core/src/main/scala/kafka/utils/ShutdownableThread.scala index 13bbc90f2a5c7..fa35aa7d94773 100644 --- a/core/src/main/scala/kafka/utils/ShutdownableThread.scala +++ b/core/src/main/scala/kafka/utils/ShutdownableThread.scala @@ -29,8 +29,13 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean private val shutdownComplete = new CountDownLatch(1) def shutdown(): Unit = { - initiateShutdown() - awaitShutdown() + if (this.isAlive) { + initiateShutdown() + awaitShutdown() + } else { + shutdownInitiated.countDown() + shutdownComplete.countDown() + } } def isShutdownComplete: Boolean = { From 292d279ab8873d1d639aa9364f6b9786f3cce76d Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Fri, 1 Feb 2019 14:11:43 -0800 Subject: [PATCH 2/9] Go through the regular shutdown logic, just avoid awaiting the shutdownComplete latch if the thread is not running. --- core/src/main/scala/kafka/utils/ShutdownableThread.scala | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/core/src/main/scala/kafka/utils/ShutdownableThread.scala b/core/src/main/scala/kafka/utils/ShutdownableThread.scala index fa35aa7d94773..ef5ed34f87bae 100644 --- a/core/src/main/scala/kafka/utils/ShutdownableThread.scala +++ b/core/src/main/scala/kafka/utils/ShutdownableThread.scala @@ -29,11 +29,10 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean private val shutdownComplete = new CountDownLatch(1) def shutdown(): Unit = { + initiateShutdown() if (this.isAlive) { - initiateShutdown() awaitShutdown() } else { - shutdownInitiated.countDown() shutdownComplete.countDown() } } From 4dfef4491c228a704be90f6cc412c1d1d5d12354 Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Wed, 6 Feb 2019 11:04:44 -0800 Subject: [PATCH 3/9] move the `Thread.isAlive` check in ShutdownableThread to awaitShutdown. --- .../main/scala/kafka/utils/ShutdownableThread.scala | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/core/src/main/scala/kafka/utils/ShutdownableThread.scala b/core/src/main/scala/kafka/utils/ShutdownableThread.scala index ef5ed34f87bae..eac889a0070b2 100644 --- a/core/src/main/scala/kafka/utils/ShutdownableThread.scala +++ b/core/src/main/scala/kafka/utils/ShutdownableThread.scala @@ -30,11 +30,7 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean def shutdown(): Unit = { initiateShutdown() - if (this.isAlive) { - awaitShutdown() - } else { - shutdownComplete.countDown() - } + awaitShutdown() } def isShutdownComplete: Boolean = { @@ -58,8 +54,10 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean * After calling initiateShutdown(), use this API to wait until the shutdown is complete */ def awaitShutdown(): Unit = { - shutdownComplete.await() - info("Shutdown completed") + if (this.isAlive) { + shutdownComplete.await() + info("Shutdown completed") + } } /** From 445a3af5094b3f55ceedc9a86677338556e88112 Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Thu, 7 Feb 2019 13:20:31 -0800 Subject: [PATCH 4/9] Allow blocking on ShutdownableThread.awaitShutdown() regardless of the thread being started or not. --- .../main/scala/kafka/utils/ShutdownableThread.scala | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/core/src/main/scala/kafka/utils/ShutdownableThread.scala b/core/src/main/scala/kafka/utils/ShutdownableThread.scala index eac889a0070b2..ef5ed34f87bae 100644 --- a/core/src/main/scala/kafka/utils/ShutdownableThread.scala +++ b/core/src/main/scala/kafka/utils/ShutdownableThread.scala @@ -30,7 +30,11 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean def shutdown(): Unit = { initiateShutdown() - awaitShutdown() + if (this.isAlive) { + awaitShutdown() + } else { + shutdownComplete.countDown() + } } def isShutdownComplete: Boolean = { @@ -54,10 +58,8 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean * After calling initiateShutdown(), use this API to wait until the shutdown is complete */ def awaitShutdown(): Unit = { - if (this.isAlive) { - shutdownComplete.await() - info("Shutdown completed") - } + shutdownComplete.await() + info("Shutdown completed") } /** From a2c0f91e57e663106164dc3e94c46c886389c8bc Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Wed, 13 Feb 2019 12:13:31 -0800 Subject: [PATCH 5/9] Don't block in `awaitShutdown()` if the thread is not started. --- core/src/main/scala/kafka/utils/ShutdownableThread.scala | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/utils/ShutdownableThread.scala b/core/src/main/scala/kafka/utils/ShutdownableThread.scala index ef5ed34f87bae..670d26e10b8b3 100644 --- a/core/src/main/scala/kafka/utils/ShutdownableThread.scala +++ b/core/src/main/scala/kafka/utils/ShutdownableThread.scala @@ -30,11 +30,7 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean def shutdown(): Unit = { initiateShutdown() - if (this.isAlive) { - awaitShutdown() - } else { - shutdownComplete.countDown() - } + awaitShutdown() } def isShutdownComplete: Boolean = { @@ -58,7 +54,8 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean * After calling initiateShutdown(), use this API to wait until the shutdown is complete */ def awaitShutdown(): Unit = { - shutdownComplete.await() + if (this.isAlive) + shutdownComplete.await() info("Shutdown completed") } From 7c6ac0d25b8c9ca48b4c190f6a11bc0ab8576054 Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Mon, 18 Feb 2019 12:16:01 -0800 Subject: [PATCH 6/9] Set ShutdownableThread started status in `start()` method to ensure it's set by the parent thread. --- .../scala/kafka/utils/ShutdownableThread.scala | 18 ++++++++++++++---- 1 file changed, 14 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/utils/ShutdownableThread.scala b/core/src/main/scala/kafka/utils/ShutdownableThread.scala index 670d26e10b8b3..0996cd22477ed 100644 --- a/core/src/main/scala/kafka/utils/ShutdownableThread.scala +++ b/core/src/main/scala/kafka/utils/ShutdownableThread.scala @@ -27,7 +27,8 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean this.logIdent = "[" + name + "]: " private val shutdownInitiated = new CountDownLatch(1) private val shutdownComplete = new CountDownLatch(1) - + @volatile private var isStarted: Boolean = false + def shutdown(): Unit = { initiateShutdown() awaitShutdown() @@ -54,9 +55,13 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean * After calling initiateShutdown(), use this API to wait until the shutdown is complete */ def awaitShutdown(): Unit = { - if (this.isAlive) - shutdownComplete.await() - info("Shutdown completed") + if (shutdownInitiated.getCount != 0) + throw new IllegalStateException("initiateShutdown() was not called before awaitShutdown()") + else { + if (isStarted) + shutdownComplete.await() + info("Shutdown completed") + } } /** @@ -76,7 +81,12 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean */ def doWork(): Unit + override def start(): Unit = { + isStarted = true + super.start() + } override def run(): Unit = { + isStarted = true info("Starting") try { while (isRunning) From 92f71e48c53593a8c4f2455dca9dfb7ca2b863b5 Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Tue, 19 Feb 2019 14:29:02 -0800 Subject: [PATCH 7/9] Remove start() override. --- core/src/main/scala/kafka/utils/ShutdownableThread.scala | 4 ---- 1 file changed, 4 deletions(-) diff --git a/core/src/main/scala/kafka/utils/ShutdownableThread.scala b/core/src/main/scala/kafka/utils/ShutdownableThread.scala index 0996cd22477ed..02d09dab96d82 100644 --- a/core/src/main/scala/kafka/utils/ShutdownableThread.scala +++ b/core/src/main/scala/kafka/utils/ShutdownableThread.scala @@ -81,10 +81,6 @@ abstract class ShutdownableThread(val name: String, val isInterruptible: Boolean */ def doWork(): Unit - override def start(): Unit = { - isStarted = true - super.start() - } override def run(): Unit = { isStarted = true info("Starting") From 4a202222000211a951f6954b68a750e7db6a9577 Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Sun, 24 Feb 2019 22:19:27 -0800 Subject: [PATCH 8/9] Adapt ControllerEventManager to the new ShutdownableThread behavior. --- .../main/scala/kafka/controller/ControllerEventManager.scala | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/controller/ControllerEventManager.scala b/core/src/main/scala/kafka/controller/ControllerEventManager.scala index 54e3a9e126b6c..63c3d26a25ce4 100644 --- a/core/src/main/scala/kafka/controller/ControllerEventManager.scala +++ b/core/src/main/scala/kafka/controller/ControllerEventManager.scala @@ -61,6 +61,7 @@ class ControllerEventManager(controllerId: Int, rateAndTimeMetrics: Map[Controll def start(): Unit = thread.start() def close(): Unit = { + thread.initiateShutdown() clearAndPut(KafkaController.ShutdownEventThread) thread.awaitShutdown() } @@ -83,7 +84,7 @@ class ControllerEventManager(controllerId: Int, rateAndTimeMetrics: Map[Controll override def doWork(): Unit = { queue.take() match { - case KafkaController.ShutdownEventThread => initiateShutdown() + case KafkaController.ShutdownEventThread => case controllerEvent => _state = controllerEvent.state From 54826d13f882d63d07a3f5f30f5329d846614ef2 Mon Sep 17 00:00:00 2001 From: Gardner Vickers Date: Mon, 25 Feb 2019 17:18:01 -0800 Subject: [PATCH 9/9] Add comment explaining why the ShutdownEventThread handler does nothing in ControllerEventManager. --- .../main/scala/kafka/controller/ControllerEventManager.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/scala/kafka/controller/ControllerEventManager.scala b/core/src/main/scala/kafka/controller/ControllerEventManager.scala index 63c3d26a25ce4..a456ce32895a9 100644 --- a/core/src/main/scala/kafka/controller/ControllerEventManager.scala +++ b/core/src/main/scala/kafka/controller/ControllerEventManager.scala @@ -84,7 +84,7 @@ class ControllerEventManager(controllerId: Int, rateAndTimeMetrics: Map[Controll override def doWork(): Unit = { queue.take() match { - case KafkaController.ShutdownEventThread => + case KafkaController.ShutdownEventThread => // The shutting down of the thread has been initiated at this point. Ignore this event. case controllerEvent => _state = controllerEvent.state