From cbfe9445386927e48841f06839a4ce86ad694363 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 16 Jan 2023 14:13:41 +0200 Subject: [PATCH] Deduplication node improvement --- .../deduplication/TbMsgDeduplicationNode.java | 128 +++++++++--------- 1 file changed, 65 insertions(+), 63 deletions(-) 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 c23d5a80fa..de04ce1d3b 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,6 +18,7 @@ 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; @@ -28,13 +29,18 @@ import org.thingsboard.rule.engine.api.TbRelationTypes; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.id.EntityId; 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.TbMsgMetaData; 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; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; @@ -60,18 +66,20 @@ public class TbMsgDeduplicationNode implements TbNode { private TbMsgDeduplicationNodeConfiguration config; - private Map> deduplicationMap; + private final Map> deduplicationMap; private long deduplicationInterval; private long lastScheduledTs; private DeduplicationId deduplicationId; + public TbMsgDeduplicationNode() { + this.deduplicationMap = new HashMap<>(); + } @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { this.config = TbNodeUtils.convert(configuration, TbMsgDeduplicationNodeConfiguration.class); this.deduplicationInterval = TimeUnit.SECONDS.toMillis(config.getInterval()); this.deduplicationId = config.getId(); - this.deduplicationMap = new HashMap<>(); scheduleTickMsg(ctx); } @@ -95,7 +103,7 @@ public class TbMsgDeduplicationNode implements TbNode { private void processOnRegularMsg(TbContext ctx, TbMsg msg) { EntityId id = getDeduplicationId(ctx, msg); - List deduplicationMsgs = deduplicationMap.computeIfAbsent(id, k -> new ArrayList<>()); + List deduplicationMsgs = deduplicationMap.computeIfAbsent(id, k -> new LinkedList<>()); if (deduplicationMsgs.size() < config.getMaxPendingMsgs()) { log.trace("[{}][{}] Adding msg: [{}][{}] to the pending msgs map ...", ctx.getSelfId(), id, msg.getId(), msg.getMetaDataTs()); deduplicationMsgs.add(msg); @@ -120,73 +128,67 @@ public class TbMsgDeduplicationNode implements TbNode { } private void processDeduplication(TbContext ctx) { - if (!deduplicationMap.isEmpty()) { - List deduplicationResults = new ArrayList<>(); - long deduplicationTimeoutMs = System.currentTimeMillis(); - deduplicationMap.forEach((entityId, tbMsgs) -> { - if (!tbMsgs.isEmpty()) { - TbMsg firstMsgInPack = getFirstPackMsg(tbMsgs); - long packStartTs = firstMsgInPack.getMetaDataTs(); - long packEndTs = packStartTs + deduplicationInterval; - boolean hasNextPack = packEndTs < deduplicationTimeoutMs; - while (hasNextPack) { - TbMsg lastMsgInPack = firstMsgInPack; - List pack = new ArrayList<>(); - pack.add(firstMsgInPack); - tbMsgs.remove(firstMsgInPack); - firstMsgInPack = null; - for (TbMsg msg : tbMsgs) { - if (msg.getMetaDataTs() > packStartTs && msg.getMetaDataTs() < packEndTs) { - pack.add(msg); - if (msg.getMetaDataTs() > lastMsgInPack.getMetaDataTs()) { - lastMsgInPack = msg; - } - } else { - if (firstMsgInPack == null || msg.getMetaDataTs() < firstMsgInPack.getMetaDataTs()) { - firstMsgInPack = msg; - } - } + if (deduplicationMap.isEmpty()) { + return; + } + List deduplicationResults = new ArrayList<>(); + long deduplicationTimeoutMs = System.currentTimeMillis(); + deduplicationMap.forEach((entityId, tbMsgs) -> { + if (tbMsgs.isEmpty()) { + return; + } + Optional> packBoundsOpt = findValidPack(tbMsgs, 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(); ) { + TbMsg msg = iterator.next(); + long msgTs = msg.getMetaDataTs(); + if (msgTs >= packBounds.getFirst() && msgTs < packBounds.getSecond()) { + pack.add(msg); + iterator.remove(); } - deduplicationResults.add(createOutMsg(entityId, pack, pack.indexOf(lastMsgInPack))); - tbMsgs.removeAll(pack); - if (firstMsgInPack == null) { - hasNextPack = false; - } else { - packStartTs = firstMsgInPack.getMetaDataTs(); - packEndTs = packStartTs + deduplicationInterval; - hasNextPack = packEndTs < deduplicationTimeoutMs; + } + deduplicationResults.add(TbMsg.newMsg( + config.getQueueName(), + config.getOutMsgType(), + entityId, + getMetadata(), + getMergedData(pack))); + } else { + TbMsg resultMsg = null; + boolean searchMin = DeduplicationStrategy.FIRST.equals(config.getStrategy()); + for (Iterator iterator = tbMsgs.iterator(); iterator.hasNext(); ) { + TbMsg msg = iterator.next(); + long msgTs = msg.getMetaDataTs(); + if (msgTs >= packBounds.getFirst() && msgTs < packBounds.getSecond()) { + iterator.remove(); + if (resultMsg == null + || (searchMin && msg.getMetaDataTs() < resultMsg.getMetaDataTs()) + || (!searchMin && msg.getMetaDataTs() > resultMsg.getMetaDataTs())) { + resultMsg = msg; + } } } + deduplicationResults.add(resultMsg); } - }); - deduplicationResults.forEach(outMsg -> enqueueForTellNextWithRetry(ctx, outMsg, 0)); - } - } - - private TbMsg getFirstPackMsg(List tbMsgs) { - TbMsg firstMsg = null; - for (TbMsg msg : tbMsgs) { - if (firstMsg == null || msg.getMetaDataTs() < firstMsg.getMetaDataTs()) { - firstMsg = msg; + packBoundsOpt = findValidPack(tbMsgs, deduplicationTimeoutMs); } - } - return firstMsg; + }); + deduplicationResults.forEach(outMsg -> enqueueForTellNextWithRetry(ctx, outMsg, 0)); } - private TbMsg createOutMsg(EntityId originator, List pack, int lastIndex) { - switch (config.getStrategy()) { - case FIRST: - return pack.get(0); - case LAST: - return pack.get(lastIndex); - default: - return TbMsg.newMsg( - config.getQueueName(), - config.getOutMsgType(), - originator, - getMetadata(), - getMergedData(pack)); - } + private Optional> findValidPack(List msgs, long deduplicationTimeoutMs) { + Optional min = msgs.stream().min(Comparator.comparing(TbMsg::getMetaDataTs)); + return min.map(minTsMsg -> { + long packStartTs = minTsMsg.getMetaDataTs(); + long packEndTs = packStartTs + deduplicationInterval; + if (packEndTs <= deduplicationTimeoutMs) { + return new TbPair<>(packStartTs, packEndTs); + } + return null; + }); } private void enqueueForTellNextWithRetry(TbContext ctx, TbMsg msg, int retryAttempt) {