From 32c59960eadb3448a8f874c619e8a242f19a05f1 Mon Sep 17 00:00:00 2001 From: harrylee Date: Thu, 23 May 2024 18:47:55 +0800 Subject: [PATCH] =?UTF-8?q?=E4=B8=BB=E9=A2=98=E8=B7=AF=E7=94=B1=E5=99=A8?= =?UTF-8?q?=EF=BC=88=E9=BB=98=E8=AE=A4&=E6=A8=A1=E5=BC=8F=E5=8C=B9?= =?UTF-8?q?=E9=85=8D=EF=BC=89=E4=BC=98=E5=8C=96=EF=BC=9A=E4=BD=BF=E7=94=A8?= =?UTF-8?q?=20ReentrantLock=20=E4=BC=98=E5=8C=96=20synchronized=20?= =?UTF-8?q?=E5=92=8C=20=E4=BD=BF=E7=94=A8java1.8=20=E7=9A=84=20removeIf=20?= =?UTF-8?q?=E4=BC=98=E5=8C=96=20remove?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../dami/bus/impl/TopicRouterDefault.java | 39 +++++++--- .../dami/bus/impl/TopicRouterPatterned.java | 43 ++++++----- .../features/demo23_unlistenall/Demo23.java | 73 +++++++++++++++++-- 3 files changed, 119 insertions(+), 36 deletions(-) diff --git a/dami/src/main/java/org/noear/dami/bus/impl/TopicRouterDefault.java b/dami/src/main/java/org/noear/dami/bus/impl/TopicRouterDefault.java index 6f88802..352572d 100644 --- a/dami/src/main/java/org/noear/dami/bus/impl/TopicRouterDefault.java +++ b/dami/src/main/java/org/noear/dami/bus/impl/TopicRouterDefault.java @@ -11,6 +11,7 @@ import org.slf4j.LoggerFactory; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.locks.ReentrantLock; /** * 主题路由器(默认啥希表实现方案) @@ -26,6 +27,8 @@ public class TopicRouterDefault implements TopicRouter { */ private final Map> pipelineMap = new LinkedHashMap<>(); + protected final ReentrantLock PIPELINE_MAP_LOCK = new ReentrantLock(); + public TopicRouterDefault() { super(); } @@ -38,10 +41,14 @@ public class TopicRouterDefault implements TopicRouter { * @param listener 监听器 */ @Override - public synchronized void add(final String topic, final int index, final TopicListener> listener) { - final TopicListenPipeline pipeline = pipelineMap.computeIfAbsent(topic, t -> new TopicListenPipeline<>()); - - pipeline.add(index, listener); + public void add(final String topic, final int index, final TopicListener> listener) { + PIPELINE_MAP_LOCK.lock(); + try { + final TopicListenPipeline pipeline = pipelineMap.computeIfAbsent(topic, t -> new TopicListenPipeline<>()); + pipeline.add(index, listener); + } finally { + PIPELINE_MAP_LOCK.unlock(); + } if (log.isDebugEnabled()) { if (MethodTopicListener.class.isAssignableFrom(listener.getClass())) { @@ -59,11 +66,18 @@ public class TopicRouterDefault implements TopicRouter { * @param listener 监听器 */ @Override - public synchronized void remove(final String topic, final TopicListener> listener) { - final TopicListenPipeline pipeline = pipelineMap.get(topic); - + public void remove(final String topic, final TopicListener> listener) { + TopicListenPipeline pipeline = pipelineMap.get(topic); if (pipeline != null) { - pipeline.remove(listener); + PIPELINE_MAP_LOCK.lock(); + try { + pipeline = pipelineMap.get(topic); + if (pipeline != null) { + pipeline.remove(listener); + } + } finally { + PIPELINE_MAP_LOCK.unlock(); + } } if (log.isDebugEnabled()) { @@ -81,8 +95,13 @@ public class TopicRouterDefault implements TopicRouter { * @param topic 主题 */ @Override - public synchronized void remove(final String topic) { - pipelineMap.remove(topic); + public void remove(final String topic) { + PIPELINE_MAP_LOCK.lock(); + try { + pipelineMap.remove(topic); + } finally { + PIPELINE_MAP_LOCK.unlock(); + } if (log.isDebugEnabled()) { log.debug("TopicRouter listener removed(@{}): all..", topic); diff --git a/dami/src/main/java/org/noear/dami/bus/impl/TopicRouterPatterned.java b/dami/src/main/java/org/noear/dami/bus/impl/TopicRouterPatterned.java index de1d526..19cf0ea 100644 --- a/dami/src/main/java/org/noear/dami/bus/impl/TopicRouterPatterned.java +++ b/dami/src/main/java/org/noear/dami/bus/impl/TopicRouterPatterned.java @@ -11,6 +11,7 @@ import org.slf4j.LoggerFactory; import java.util.ArrayList; import java.util.Comparator; import java.util.List; +import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; @@ -29,6 +30,9 @@ public class TopicRouterPatterned implements TopicRouter { */ private final List> routingList = new ArrayList<>(); + protected final ReentrantLock ROUTING_LIST_LOCK = new ReentrantLock(); + + /** * 路由工厂 */ @@ -47,9 +51,13 @@ public class TopicRouterPatterned implements TopicRouter { * @param listener 监听器 */ @Override - public synchronized void add(final String topic, final int index, final TopicListener> listener) { - routingList.add(routerFactory.create(topic, index, listener)); - + public void add(final String topic, final int index, final TopicListener> listener) { + ROUTING_LIST_LOCK.lock(); + try { + routingList.add(routerFactory.create(topic, index, listener)); + } finally { + ROUTING_LIST_LOCK.unlock(); + } if (log.isDebugEnabled()) { if (MethodTopicListener.class.isAssignableFrom(listener.getClass())) { log.debug("TopicRouter listener added(@{}): {}", topic, listener); @@ -66,17 +74,15 @@ public class TopicRouterPatterned implements TopicRouter { * @param listener 监听器 */ @Override - public synchronized void remove(final String topic, final TopicListener> listener) { - for (int i = 0; i < routingList.size(); i++) { - Routing routing = routingList.get(i); - if (routing.matches(topic)) { - if (routing.getListener() == listener) { - routingList.remove(i); - i--; - } - } + public void remove(final String topic, final TopicListener> listener) { + ROUTING_LIST_LOCK.lock(); + try { + routingList.removeIf(routing -> + routing.matches(topic) + && routing.getListener() == listener); + } finally { + ROUTING_LIST_LOCK.unlock(); } - if (log.isDebugEnabled()) { if (MethodTopicListener.class.isAssignableFrom(listener.getClass())) { log.debug("TopicRouter listener removed(@{}): {}", topic, listener); @@ -93,12 +99,11 @@ public class TopicRouterPatterned implements TopicRouter { */ @Override public void remove(String topic) { - for (int i = 0; i < routingList.size(); i++) { - Routing routing = routingList.get(i); - if (routing.matches(topic)) { - routingList.remove(i); - i--; - } + ROUTING_LIST_LOCK.lock(); + try { + routingList.removeIf(routing -> routing.matches(topic)); + } finally { + ROUTING_LIST_LOCK.unlock(); } if (log.isDebugEnabled()) { diff --git a/dami/src/test/java/features/demo23_unlistenall/Demo23.java b/dami/src/test/java/features/demo23_unlistenall/Demo23.java index e59c30a..bead26b 100644 --- a/dami/src/test/java/features/demo23_unlistenall/Demo23.java +++ b/dami/src/test/java/features/demo23_unlistenall/Demo23.java @@ -6,6 +6,10 @@ import org.noear.dami.bus.DamiBusImpl; import org.noear.dami.bus.Payload; import org.noear.dami.bus.TopicListener; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + public class Demo23 { static String topic = "demo.hello"; //定义实例,避免单测干扰 //开发时用:Dami.bus() @@ -13,28 +17,83 @@ public class Demo23 { @Test public void main() throws InterruptedException { - TopicListener> aListener = payload -> { + final TopicListener> aListener = payload -> { System.out.println("i'm a:" + payload); }; - TopicListener> bListener = payload -> { + final TopicListener> bListener = payload -> { System.out.println("i'm b:" + payload); }; - TopicListener> cListener = payload -> { + final TopicListener> cListener = payload -> { System.out.println("i'm b:" + payload); }; //监听事件 - bus.listen(topic,aListener); - bus.listen(topic,bListener); - bus.listen(topic,cListener); + bus.listen(topic, aListener); + bus.listen(topic, bListener); + bus.listen(topic, cListener); System.out.println("------------------ 从未移除监听器之前 ------------------"); //发送事件 bus.send(topic, "aaaaaa"); System.out.println("------------------ 移除 a 监听器之后 ------------------"); - bus.unlisten(topic,aListener); + bus.unlisten(topic, aListener); bus.send(topic, "bbbbbb"); System.out.println("------------------ 移除所有监听器之后 ------------------"); bus.unlisten(topic); bus.send(topic, "cccccc"); + + + // 验证多线程情况下 + CountDownLatch listen = new CountDownLatch(3); + CompletableFuture.runAsync(()->{ + bus.listen(topic, aListener); + listen.countDown(); + }); + CompletableFuture.runAsync(()->{ + bus.listen(topic, bListener); + listen.countDown(); + }); + CompletableFuture.runAsync(()->{ + bus.listen(topic, cListener); + listen.countDown(); + }); + assert listen.await(1, TimeUnit.SECONDS); + + System.out.println("------------------ runAsync 从未移除监听器之前 ------------------"); + //发送事件 + bus.send(topic, "runAsync aaaaaa"); + System.out.println("------------------ runAsync 移除 a 监听器之后 ------------------"); + CountDownLatch unlisten = new CountDownLatch(2); + CompletableFuture.runAsync(()->{ + bus.unlisten(topic, aListener); + unlisten.countDown(); + }); + CompletableFuture.runAsync(()->{ + // 等待 + try { + Thread.sleep(1000); + } catch (InterruptedException e) { + } + bus.send(topic, "runAsync bbbbbb"); + unlisten.countDown(); + }); + assert unlisten.await(2, TimeUnit.SECONDS); + + CountDownLatch unlistenall = new CountDownLatch(2); + + System.out.println("------------------ runAsync 移除所有监听器之后 ------------------"); + CompletableFuture.runAsync(()->{ + bus.unlisten(topic); + unlistenall.countDown(); + }); + CompletableFuture.runAsync(()->{ + // 等待 + try { + Thread.sleep(1000); + } catch (InterruptedException e) { + } + bus.send(topic, "runAsync cccccc"); + unlistenall.countDown(); + }); + assert unlistenall.await(2, TimeUnit.SECONDS); } } -- Gitee