Browse Source

Deduplication node improvement

feature/deduplication-node
Andrii Shvaika 4 years ago
parent
commit
cbfe944538
  1. 128
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java

128
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<EntityId, List<TbMsg>> deduplicationMap;
private final Map<EntityId, List<TbMsg>> 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<TbMsg> deduplicationMsgs = deduplicationMap.computeIfAbsent(id, k -> new ArrayList<>());
List<TbMsg> 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<TbMsg> 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<TbMsg> 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<TbMsg> deduplicationResults = new ArrayList<>();
long deduplicationTimeoutMs = System.currentTimeMillis();
deduplicationMap.forEach((entityId, tbMsgs) -> {
if (tbMsgs.isEmpty()) {
return;
}
Optional<TbPair<Long, Long>> packBoundsOpt = findValidPack(tbMsgs, 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(); ) {
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<TbMsg> 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<TbMsg> 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<TbMsg> 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<TbPair<Long, Long>> findValidPack(List<TbMsg> msgs, long deduplicationTimeoutMs) {
Optional<TbMsg> 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) {

Loading…
Cancel
Save