@ -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 ( ) ,
entity Id,
deduplication Id,
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 ) {