diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java index 133b1d00cb..af7412bfe8 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java @@ -16,13 +16,14 @@ package org.thingsboard.rule.engine.math; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.github.benmanes.caffeine.cache.Caffeine; +import com.github.benmanes.caffeine.cache.LoadingCache; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.SettableFuture; import net.objecthunter.exp4j.Expression; import net.objecthunter.exp4j.ExpressionBuilder; -import org.springframework.util.ConcurrentReferenceHashMap; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.rule.engine.api.RuleNode; @@ -45,9 +46,9 @@ import org.thingsboard.server.common.msg.TbMsg; import java.math.BigDecimal; import java.math.RoundingMode; +import java.time.Duration; import java.util.List; import java.util.Optional; -import java.util.concurrent.ConcurrentMap; import java.util.function.BiFunction; import java.util.function.Function; import java.util.stream.Collectors; @@ -80,7 +81,10 @@ import static org.thingsboard.rule.engine.math.TbMathArgumentType.CONSTANT; ) public class TbMathNode implements TbNode { - private static final ConcurrentMap locks = new ConcurrentReferenceHashMap<>(16, ConcurrentReferenceHashMap.ReferenceType.WEAK); + private static final LoadingCache queues = Caffeine.newBuilder() + .expireAfterAccess(Duration.ofDays(1L)) + .build(SemaphoreWithTbMsgQueue::new); + private final ThreadLocal customExpression = new ThreadLocal<>(); private TbMathNodeConfiguration config; private boolean msgBodyToJsonConversionRequired; @@ -106,8 +110,8 @@ public class TbMathNode implements TbNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { - locks.computeIfAbsent(msg.getOriginator(), SemaphoreWithTbMsgQueue::new) - .addToQueueAndTryProcess(msg, ctx, this::processMsgAsync); + var processingQueue = queues.get(msg.getOriginator()); + processingQueue.addToQueueAndTryProcess(msg, ctx, this::processMsgAsync); } ListenableFuture processMsgAsync(TbContext ctx, TbMsg msg) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java index bdd6d0215c..33601952d2 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java @@ -17,6 +17,8 @@ package org.thingsboard.rule.engine.metadata; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.github.benmanes.caffeine.cache.Caffeine; +import com.github.benmanes.caffeine.cache.LoadingCache; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; @@ -37,9 +39,11 @@ import org.thingsboard.server.common.data.msg.TbNodeConnectionType; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; import java.math.BigDecimal; import java.math.RoundingMode; +import java.time.Duration; import java.util.Map; @RuleNode( @@ -58,7 +62,7 @@ import java.util.Map; public class CalculateDeltaNode implements TbNode { private Map cache; - private Map locks; + private LoadingCache queues; private CalculateDeltaNodeConfiguration config; @@ -74,7 +78,9 @@ public class CalculateDeltaNode implements TbNode { if (config.isAddPeriodBetweenMsgs() && StringUtils.isBlank(config.getPeriodValueKey())) { throw new TbNodeException("Period value key should be specified!", true); } - locks = new ConcurrentReferenceHashMap<>(16, ConcurrentReferenceHashMap.ReferenceType.WEAK); + queues = Caffeine.newBuilder() + .expireAfterAccess(Duration.ofDays(1L)) + .build(SemaphoreWithTbMsgQueue::new); if (config.isUseCache()) { cache = new ConcurrentReferenceHashMap<>(16, ConcurrentReferenceHashMap.ReferenceType.SOFT); } @@ -91,13 +97,18 @@ public class CalculateDeltaNode implements TbNode { ctx.tellNext(msg, TbNodeConnectionType.OTHER); return; } - locks.computeIfAbsent(msg.getOriginator(), SemaphoreWithTbMsgQueue::new) - .addToQueueAndTryProcess(msg, ctx, this::processMsgAsync); + var processingQueue = queues.get(msg.getOriginator()); + processingQueue.addToQueueAndTryProcess(msg, ctx, this::processMsgAsync); + } + + @Override + public void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) { + queues.asMap().keySet().removeIf(entityId -> !ctx.isLocalEntity(entityId)); } @Override public void destroy() { - locks.clear(); + queues.invalidateAll(); if (config.isUseCache()) { cache.clear(); }