diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/DeduplicationData.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/DeduplicationData.java new file mode 100644 index 0000000000..58f170e0ab --- /dev/null +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/DeduplicationData.java @@ -0,0 +1,45 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.rule.engine.deduplication; + +import lombok.Data; +import org.thingsboard.server.common.msg.TbMsg; + +import java.util.LinkedList; +import java.util.List; + +@Data +public class DeduplicationData { + + private final List msgList; + private boolean tickScheduled; + + public DeduplicationData() { + msgList = new LinkedList<>(); + } + + public int size() { + return msgList.size(); + } + + public void add(TbMsg msg) { + msgList.add(msg); + } + + public boolean isEmpty() { + return msgList.isEmpty(); + } +} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java index de04ce1d3b..b63894a61e 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java @@ -18,7 +18,6 @@ package org.thingsboard.rule.engine.deduplication; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; -import org.springframework.data.util.Pair; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; @@ -37,7 +36,6 @@ import java.util.ArrayList; import java.util.Comparator; import java.util.HashMap; import java.util.Iterator; -import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Optional; @@ -61,14 +59,14 @@ import java.util.concurrent.TimeUnit; public class TbMsgDeduplicationNode implements TbNode { private static final String TB_MSG_DEDUPLICATION_TIMEOUT_MSG = "TbMsgDeduplicationNodeMsg"; - private static final int TB_MSG_DEDUPLICATION_TIMEOUT = 5000; public static final int TB_MSG_DEDUPLICATION_RETRY_DELAY = 10; + private static final String EMPTY_DATA = ""; + private static final TbMsgMetaData EMPTY_META_DATA = new TbMsgMetaData(); private TbMsgDeduplicationNodeConfiguration config; - private final Map> deduplicationMap; + private final Map deduplicationMap; private long deduplicationInterval; - private long lastScheduledTs; private DeduplicationId deduplicationId; public TbMsgDeduplicationNode() { @@ -80,17 +78,12 @@ public class TbMsgDeduplicationNode implements TbNode { this.config = TbNodeUtils.convert(configuration, TbMsgDeduplicationNodeConfiguration.class); this.deduplicationInterval = TimeUnit.SECONDS.toMillis(config.getInterval()); this.deduplicationId = config.getId(); - scheduleTickMsg(ctx); } @Override public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException { if (TB_MSG_DEDUPLICATION_TIMEOUT_MSG.equals(msg.getType())) { - try { - processDeduplication(ctx); - } finally { - scheduleTickMsg(ctx); - } + processDeduplication(ctx, msg.getOriginator()); } else { processOnRegularMsg(ctx, msg); } @@ -103,11 +96,12 @@ public class TbMsgDeduplicationNode implements TbNode { private void processOnRegularMsg(TbContext ctx, TbMsg msg) { EntityId id = getDeduplicationId(ctx, msg); - List deduplicationMsgs = deduplicationMap.computeIfAbsent(id, k -> new LinkedList<>()); + DeduplicationData deduplicationMsgs = deduplicationMap.computeIfAbsent(id, k -> new DeduplicationData()); if (deduplicationMsgs.size() < config.getMaxPendingMsgs()) { log.trace("[{}][{}] Adding msg: [{}][{}] to the pending msgs map ...", ctx.getSelfId(), id, msg.getId(), msg.getMetaDataTs()); deduplicationMsgs.add(msg); ctx.ack(msg); + scheduleTickMsg(ctx, id, deduplicationMsgs); } else { log.trace("[{}] Max limit of pending messages reached for deduplication id: [{}]", ctx.getSelfId(), id); ctx.tellFailure(msg, new RuntimeException("[" + ctx.getSelfId() + "] Max limit of pending messages reached for deduplication id: [" + id + "]")); @@ -127,22 +121,25 @@ public class TbMsgDeduplicationNode implements TbNode { } } - private void processDeduplication(TbContext ctx) { - if (deduplicationMap.isEmpty()) { + private void processDeduplication(TbContext ctx, EntityId deduplicationId) { + DeduplicationData data = deduplicationMap.get(deduplicationId); + if (data == null) { + return; + } + data.setTickScheduled(false); + if (data.isEmpty()) { return; } - List deduplicationResults = new ArrayList<>(); long deduplicationTimeoutMs = System.currentTimeMillis(); - deduplicationMap.forEach((entityId, tbMsgs) -> { - if (tbMsgs.isEmpty()) { - return; - } - Optional> packBoundsOpt = findValidPack(tbMsgs, deduplicationTimeoutMs); + try { + List deduplicationResults = new ArrayList<>(); + List msgList = data.getMsgList(); + Optional> packBoundsOpt = findValidPack(msgList, deduplicationTimeoutMs); while (packBoundsOpt.isPresent()) { TbPair packBounds = packBoundsOpt.get(); if (DeduplicationStrategy.ALL.equals(config.getStrategy())) { List pack = new ArrayList<>(); - for (Iterator iterator = tbMsgs.iterator(); iterator.hasNext(); ) { + for (Iterator iterator = msgList.iterator(); iterator.hasNext(); ) { TbMsg msg = iterator.next(); long msgTs = msg.getMetaDataTs(); if (msgTs >= packBounds.getFirst() && msgTs < packBounds.getSecond()) { @@ -153,13 +150,13 @@ public class TbMsgDeduplicationNode implements TbNode { deduplicationResults.add(TbMsg.newMsg( config.getQueueName(), config.getOutMsgType(), - entityId, + deduplicationId, getMetadata(), getMergedData(pack))); } else { TbMsg resultMsg = null; boolean searchMin = DeduplicationStrategy.FIRST.equals(config.getStrategy()); - for (Iterator iterator = tbMsgs.iterator(); iterator.hasNext(); ) { + for (Iterator iterator = msgList.iterator(); iterator.hasNext(); ) { TbMsg msg = iterator.next(); long msgTs = msg.getMetaDataTs(); if (msgTs >= packBounds.getFirst() && msgTs < packBounds.getSecond()) { @@ -173,10 +170,21 @@ public class TbMsgDeduplicationNode implements TbNode { } deduplicationResults.add(resultMsg); } - packBoundsOpt = findValidPack(tbMsgs, deduplicationTimeoutMs); + packBoundsOpt = findValidPack(msgList, deduplicationTimeoutMs); } - }); - deduplicationResults.forEach(outMsg -> enqueueForTellNextWithRetry(ctx, outMsg, 0)); + deduplicationResults.forEach(outMsg -> enqueueForTellNextWithRetry(ctx, outMsg, 0)); + } finally { + if (!data.isEmpty()) { + scheduleTickMsg(ctx, deduplicationId, data); + } + } + } + + private void scheduleTickMsg(TbContext ctx, EntityId deduplicationId, DeduplicationData data) { + if (!data.isTickScheduled()) { + scheduleTickMsg(ctx, deduplicationId); + data.setTickScheduled(true); + } } private Optional> findValidPack(List msgs, long deduplicationTimeoutMs) { @@ -206,15 +214,8 @@ public class TbMsgDeduplicationNode implements TbNode { } } - private void scheduleTickMsg(TbContext ctx) { - long curTs = System.currentTimeMillis(); - if (lastScheduledTs == 0L) { - lastScheduledTs = curTs; - } - lastScheduledTs += TB_MSG_DEDUPLICATION_TIMEOUT; - long curDelay = Math.max(0L, (lastScheduledTs - curTs)); - TbMsg tickMsg = ctx.newMsg(null, TB_MSG_DEDUPLICATION_TIMEOUT_MSG, ctx.getSelfId(), new TbMsgMetaData(), ""); - ctx.tellSelf(tickMsg, curDelay); + private void scheduleTickMsg(TbContext ctx, EntityId deduplicationId) { + ctx.tellSelf(ctx.newMsg(null, TB_MSG_DEDUPLICATION_TIMEOUT_MSG, deduplicationId, EMPTY_META_DATA, EMPTY_DATA), deduplicationInterval + 1); } private String getMergedData(List msgs) {