diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java index 1985c22d6e..824ba48fc6 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java @@ -71,7 +71,7 @@ public class AckMessageProcessor implements NettyRequestProcessor { public void shutdownPopReviveService() { for (PopReviveService popReviveService : popReviveServices) { - popReviveService.stop(); + popReviveService.shutdown(); } } diff --git a/common/src/main/java/org/apache/rocketmq/common/ServiceThread.java b/common/src/main/java/org/apache/rocketmq/common/ServiceThread.java index 4b7da90df5..95dc8b9800 100644 --- a/common/src/main/java/org/apache/rocketmq/common/ServiceThread.java +++ b/common/src/main/java/org/apache/rocketmq/common/ServiceThread.java @@ -67,9 +67,8 @@ public abstract class ServiceThread implements Runnable { this.stopped = true; log.info("shutdown thread[{}] interrupt={} ", getServiceName(), interrupt); - if (hasNotified.compareAndSet(false, true)) { - waitPoint.countDown(); // notify - } + //if thead is waiting, wakeup it + wakeup(); try { if (interrupt) { @@ -91,28 +90,6 @@ public abstract class ServiceThread implements Runnable { return JOIN_TIME; } - @Deprecated - public void stop() { - this.stop(false); - } - - @Deprecated - public void stop(final boolean interrupt) { - if (!started.get()) { - return; - } - this.stopped = true; - log.info("stop thread[{}],interrupt={} ", this.getServiceName(), interrupt); - - if (hasNotified.compareAndSet(false, true)) { - waitPoint.countDown(); // notify - } - - if (interrupt) { - this.thread.interrupt(); - } - } - public void makeStop() { if (!started.get()) { return; diff --git a/common/src/test/java/org/apache/rocketmq/common/ServiceThreadTest.java b/common/src/test/java/org/apache/rocketmq/common/ServiceThreadTest.java index 24a4af89bc..93208bcb7f 100644 --- a/common/src/test/java/org/apache/rocketmq/common/ServiceThreadTest.java +++ b/common/src/test/java/org/apache/rocketmq/common/ServiceThreadTest.java @@ -31,12 +31,6 @@ public class ServiceThreadTest { shutdown(true, true); } - @Test - public void testStop() { - stop(true); - stop(false); - } - @Test public void testMakeStop() { ServiceThread testServiceThread = startTestServiceThread(); @@ -116,23 +110,4 @@ public class ServiceThreadTest { assertEquals(true, testServiceThread.hasNotified.get()); assertEquals(0, testServiceThread.waitPoint.getCount()); } - - public void stop(boolean interrupt) { - ServiceThread testServiceThread = startTestServiceThread(); - stop0(interrupt, testServiceThread); - // repeat - stop0(interrupt, testServiceThread); - } - - private void stop0(boolean interrupt, ServiceThread testServiceThread) { - if (interrupt) { - testServiceThread.stop(true); - } else { - testServiceThread.stop(); - } - assertEquals(true, testServiceThread.isStopped()); - assertEquals(true, testServiceThread.hasNotified.get()); - assertEquals(0, testServiceThread.waitPoint.getCount()); - } - }