From 46a58ca82bb17c4434de39c52e34203dfd8fd417 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 3 Jun 2025 12:47:55 +0300 Subject: [PATCH 1/4] Edqs - VersionStore - Use local cache instead of caffeine to reduce memory heap size --- .../server/edqs/processor/EdqsProcessor.java | 1 + .../server/edqs/util/VersionsStore.java | 47 +++++++++++++++---- 2 files changed, 39 insertions(+), 9 deletions(-) diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java index 510d2c3a41..0e74cb98fa 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java @@ -277,6 +277,7 @@ public class EdqsProcessor implements TbQueueHandler, eventConsumer.awaitStop(); responseTemplate.stop(); stateService.stop(); + versionsStore.shutdown(); } } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java index ba3263eec2..9d4c67c4c2 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java @@ -15,31 +15,35 @@ */ package org.thingsboard.server.edqs.util; -import com.github.benmanes.caffeine.cache.Cache; -import com.github.benmanes.caffeine.cache.Caffeine; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.edqs.EdqsObjectKey; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @Slf4j public class VersionsStore { - private final Cache versions; + private final ConcurrentMap> versions = new ConcurrentHashMap<>(); + private final long expirationMillis; + private final ScheduledExecutorService cleaner = Executors.newSingleThreadScheduledExecutor(); public VersionsStore(int ttlMinutes) { - this.versions = Caffeine.newBuilder() - .expireAfterWrite(ttlMinutes, TimeUnit.MINUTES) - .build(); + this.expirationMillis = TimeUnit.MINUTES.toMillis(ttlMinutes); + startCleanupTask(); } public boolean isNew(EdqsObjectKey key, Long version) { AtomicBoolean isNew = new AtomicBoolean(false); - versions.asMap().compute(key, (k, prevVersion) -> { - if (prevVersion == null || prevVersion <= version) { + versions.compute(key, (k, prevVersion) -> { + if (prevVersion == null || prevVersion.value <= version) { isNew.set(true); - return version; + return new TimedValue<>(version); } else { log.debug("[{}] Version {} is outdated, the latest is {}", key, version, prevVersion); return prevVersion; @@ -48,4 +52,29 @@ public class VersionsStore { return isNew.get(); } + private void startCleanupTask() { + cleaner.scheduleAtFixedRate(() -> { + long now = System.currentTimeMillis(); + for (Map.Entry> entry : versions.entrySet()) { + if (now - entry.getValue().lastUpdated > expirationMillis) { + versions.remove(entry.getKey(), entry.getValue()); + } + } + }, expirationMillis, expirationMillis, TimeUnit.MILLISECONDS); + } + + public void shutdown() { + cleaner.shutdown(); + } + + private static class TimedValue { + private final long lastUpdated; + private final V value; + + public TimedValue(V value) { + this.value = value; + this.lastUpdated = System.currentTimeMillis(); + } + } + } From ccdcbc635043bfb05628612d0179e03cfd9bfb18 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 3 Jun 2025 15:47:19 +0300 Subject: [PATCH 2/4] VersionsStore - use long intead of Long to decrease heap size --- .../thingsboard/server/edqs/util/VersionsStore.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java index 9d4c67c4c2..c8c8f76761 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java @@ -29,7 +29,7 @@ import java.util.concurrent.atomic.AtomicBoolean; @Slf4j public class VersionsStore { - private final ConcurrentMap> versions = new ConcurrentHashMap<>(); + private final ConcurrentMap versions = new ConcurrentHashMap<>(); private final long expirationMillis; private final ScheduledExecutorService cleaner = Executors.newSingleThreadScheduledExecutor(); @@ -43,7 +43,7 @@ public class VersionsStore { versions.compute(key, (k, prevVersion) -> { if (prevVersion == null || prevVersion.value <= version) { isNew.set(true); - return new TimedValue<>(version); + return new TimedValue(version); } else { log.debug("[{}] Version {} is outdated, the latest is {}", key, version, prevVersion); return prevVersion; @@ -55,7 +55,7 @@ public class VersionsStore { private void startCleanupTask() { cleaner.scheduleAtFixedRate(() -> { long now = System.currentTimeMillis(); - for (Map.Entry> entry : versions.entrySet()) { + for (Map.Entry entry : versions.entrySet()) { if (now - entry.getValue().lastUpdated > expirationMillis) { versions.remove(entry.getKey(), entry.getValue()); } @@ -67,11 +67,11 @@ public class VersionsStore { cleaner.shutdown(); } - private static class TimedValue { + private static class TimedValue { private final long lastUpdated; - private final V value; + private final long value; - public TimedValue(V value) { + public TimedValue(long value) { this.value = value; this.lastUpdated = System.currentTimeMillis(); } From 1d5c4ac7ab5f978fb05339cc6e897d40c2e8fbc1 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 3 Jun 2025 15:48:26 +0300 Subject: [PATCH 3/4] VersionsStore - added try/catch for cleanup task --- .../thingsboard/server/edqs/util/VersionsStore.java | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java index c8c8f76761..f348e9cf9e 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/util/VersionsStore.java @@ -54,11 +54,15 @@ public class VersionsStore { private void startCleanupTask() { cleaner.scheduleAtFixedRate(() -> { - long now = System.currentTimeMillis(); - for (Map.Entry entry : versions.entrySet()) { - if (now - entry.getValue().lastUpdated > expirationMillis) { - versions.remove(entry.getKey(), entry.getValue()); + try { + long now = System.currentTimeMillis(); + for (Map.Entry entry : versions.entrySet()) { + if (now - entry.getValue().lastUpdated > expirationMillis) { + versions.remove(entry.getKey(), entry.getValue()); + } } + } catch (Exception e) { + log.error("Cleanup task failed", e); } }, expirationMillis, expirationMillis, TimeUnit.MILLISECONDS); } From 36a2b3f66624702adb12df9b5911e6536043090f Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 4 Jun 2025 13:13:29 +0300 Subject: [PATCH 4/4] KafkaEdqsStateService - added versionsStore.shutdown() --- .../org/thingsboard/server/edqs/state/KafkaEdqsStateService.java | 1 + 1 file changed, 1 insertion(+) diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java index 66bbb7a68a..7e2e99e662 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java @@ -224,6 +224,7 @@ public class KafkaEdqsStateService implements EdqsStateService { stateConsumer.awaitStop(); eventsToBackupConsumer.stop(); stateProducer.stop(); + versionsStore.shutdown(); } }