|
|
|
@ -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<EntityId, List<TbMsg>> deduplicationMap; |
|
|
|
private final Map<EntityId, DeduplicationData> 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<TbMsg> 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<TbMsg> deduplicationResults = new ArrayList<>(); |
|
|
|
long deduplicationTimeoutMs = System.currentTimeMillis(); |
|
|
|
deduplicationMap.forEach((entityId, tbMsgs) -> { |
|
|
|
if (tbMsgs.isEmpty()) { |
|
|
|
return; |
|
|
|
} |
|
|
|
Optional<TbPair<Long, Long>> packBoundsOpt = findValidPack(tbMsgs, deduplicationTimeoutMs); |
|
|
|
try { |
|
|
|
List<TbMsg> deduplicationResults = new ArrayList<>(); |
|
|
|
List<TbMsg> msgList = data.getMsgList(); |
|
|
|
Optional<TbPair<Long, Long>> packBoundsOpt = findValidPack(msgList, deduplicationTimeoutMs); |
|
|
|
while (packBoundsOpt.isPresent()) { |
|
|
|
TbPair<Long, Long> packBounds = packBoundsOpt.get(); |
|
|
|
if (DeduplicationStrategy.ALL.equals(config.getStrategy())) { |
|
|
|
List<TbMsg> pack = new ArrayList<>(); |
|
|
|
for (Iterator<TbMsg> iterator = tbMsgs.iterator(); iterator.hasNext(); ) { |
|
|
|
for (Iterator<TbMsg> 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<TbMsg> iterator = tbMsgs.iterator(); iterator.hasNext(); ) { |
|
|
|
for (Iterator<TbMsg> 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<TbPair<Long, Long>> findValidPack(List<TbMsg> 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<TbMsg> msgs) { |
|
|
|
|