From 37940867263ed7cc8bca68187a68e3401460b815 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 28 Aug 2020 15:24:31 +0300 Subject: [PATCH 1/2] added logs for in memory queue --- .../server/queue/memory/InMemoryStorage.java | 18 ++++++++++++++++++ .../queue/memory/InMemoryTbQueueProducer.java | 2 +- 2 files changed, 19 insertions(+), 1 deletion(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java index f51f5b8ae0..c4414089c4 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java @@ -23,16 +23,29 @@ import java.util.Collections; import java.util.List; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @Slf4j public final class InMemoryStorage { private static InMemoryStorage instance; private final ConcurrentHashMap> storage; + private static ScheduledExecutorService statExecutor; private InMemoryStorage() { storage = new ConcurrentHashMap<>(); + statExecutor = Executors.newSingleThreadScheduledExecutor(); + statExecutor.scheduleAtFixedRate(this::printStats, 30, 30, TimeUnit.SECONDS); + } + + private void printStats() { + storage.forEach((topic, queue) -> { + if (queue.size() > 0) { + log.debug("Topic: [{}], Queue size: [{}]", topic, queue.size()); + } + }); } public static InMemoryStorage getInstance() { @@ -77,4 +90,9 @@ public final class InMemoryStorage { storage.clear(); } + public void destroy() { + if (statExecutor != null) { + statExecutor.shutdownNow(); + } + } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java index cfcd788a16..84a9a1fdf0 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java @@ -53,6 +53,6 @@ public class InMemoryTbQueueProducer implements TbQueuePro @Override public void stop() { - + storage.destroy(); } } From 17ee07e09a0ae11b18e9aeb280cbd7a79afe3e8b Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 28 Aug 2020 16:01:43 +0300 Subject: [PATCH 2/2] changed log time --- .../org/thingsboard/server/queue/memory/InMemoryStorage.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java index c4414089c4..994ca26305 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java @@ -37,7 +37,7 @@ public final class InMemoryStorage { private InMemoryStorage() { storage = new ConcurrentHashMap<>(); statExecutor = Executors.newSingleThreadScheduledExecutor(); - statExecutor.scheduleAtFixedRate(this::printStats, 30, 30, TimeUnit.SECONDS); + statExecutor.scheduleAtFixedRate(this::printStats, 60, 60, TimeUnit.SECONDS); } private void printStats() {