From 8914cdb70784fec737f49a23c1a8437739fb9985 Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Mon, 1 Nov 2021 13:23:38 +0200 Subject: [PATCH 1/3] fixed the reprocessing logic for rule engine messages when the strategy is not set to reprocess them --- .../TbRuleEngineProcessingStrategyFactory.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java index f29377445d..f2ee70f333 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.queue.processing; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; +import org.springframework.util.CollectionUtils; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.queue.TbMsgCallback; import org.thingsboard.server.gen.transport.TransportProtos; @@ -77,6 +78,7 @@ public class TbRuleEngineProcessingStrategyFactory { @Override public TbRuleEngineProcessingDecision analyze(TbRuleEngineProcessingResult result) { if (result.isSuccess()) { + log.debug("[{}] The result of the msg pack processing is successful, going to proceed with processing of the following msgs", queueName); return new TbRuleEngineProcessingDecision(true, null); } else { if (retryCount == 0) { @@ -91,6 +93,7 @@ public class TbRuleEngineProcessingStrategyFactory { log.debug("[{}] Skip reprocess of the rule engine pack due to max allowed failure percentage", queueName); return new TbRuleEngineProcessingDecision(true, null); } else { + log.debug("[{}] The result of msg pack processing is unsuccessful, checking unprocessed msgs and going to reprocess them", queueName); ConcurrentMap> toReprocess = new ConcurrentHashMap<>(initialTotalCount); if (retryFailed) { result.getFailedMap().forEach(toReprocess::put); @@ -101,6 +104,15 @@ public class TbRuleEngineProcessingStrategyFactory { if (retrySuccessful) { result.getSuccessMap().forEach(toReprocess::put); } + if (CollectionUtils.isEmpty(toReprocess)) { + log.debug("[{}] Stopping the reprocessing logic due to reprocessing map is empty", queueName); + if (retryFailed) { + log.debug("[{}] Probably there are timedOut msgs, but the strategy is not set to reprocess them. Reprocess of the failed msgs is set", queueName); + } else if (retryTimeout) { + log.debug("[{}] Probably there are failed msgs, but the strategy is not set to reprocess them. Reprocess of the timedOut msgs is set", queueName); + } + return new TbRuleEngineProcessingDecision(true, null); + } log.debug("[{}] Going to reprocess {} messages", queueName, toReprocess.size()); if (log.isTraceEnabled()) { toReprocess.forEach((id, msg) -> log.trace("Going to reprocess [{}]: {}", id, TbMsg.fromBytes(result.getQueueName(), msg.getValue().getTbMsg().toByteArray(), TbMsgCallback.EMPTY))); From c49297b4335c14361637e710d64e3a55aeef0354 Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Mon, 1 Nov 2021 14:57:16 +0200 Subject: [PATCH 2/3] logging improvements --- .../TbRuleEngineProcessingStrategyFactory.java | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java index f2ee70f333..2749dd7fc5 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java @@ -97,19 +97,22 @@ public class TbRuleEngineProcessingStrategyFactory { ConcurrentMap> toReprocess = new ConcurrentHashMap<>(initialTotalCount); if (retryFailed) { result.getFailedMap().forEach(toReprocess::put); + } else if (log.isDebugEnabled() && !result.getFailedMap().isEmpty()) { + log.debug("[{}] Skipped {} failed messages due to the processing strategy configuration", queueName, result.getFailedMap().size()); } if (retryTimeout) { result.getPendingMap().forEach(toReprocess::put); + } else if (log.isDebugEnabled() && !result.getPendingMap().isEmpty()) { + log.debug("[{}] Skipped {} timedOut messages due to the processing strategy configuration", queueName, result.getPendingMap().size()); } if (retrySuccessful) { result.getSuccessMap().forEach(toReprocess::put); + } else if (log.isDebugEnabled() && !result.getSuccessMap().isEmpty()) { + log.debug("[{}] Skipped {} successful messages due to the processing strategy configuration", queueName, result.getSuccessMap().size()); } if (CollectionUtils.isEmpty(toReprocess)) { - log.debug("[{}] Stopping the reprocessing logic due to reprocessing map is empty", queueName); - if (retryFailed) { - log.debug("[{}] Probably there are timedOut msgs, but the strategy is not set to reprocess them. Reprocess of the failed msgs is set", queueName); - } else if (retryTimeout) { - log.debug("[{}] Probably there are failed msgs, but the strategy is not set to reprocess them. Reprocess of the timedOut msgs is set", queueName); + if (log.isDebugEnabled()) { + log.debug("[{}] Stopping the reprocessing logic due to reprocessing map is empty", queueName); } return new TbRuleEngineProcessingDecision(true, null); } From d78c1b2de42c98dbddf642cd58f3383f71f59a82 Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Thu, 4 Nov 2021 13:15:13 +0200 Subject: [PATCH 3/3] Update TbRuleEngineProcessingStrategyFactory.java --- .../processing/TbRuleEngineProcessingStrategyFactory.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java index 2749dd7fc5..a94731696f 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/TbRuleEngineProcessingStrategyFactory.java @@ -78,7 +78,7 @@ public class TbRuleEngineProcessingStrategyFactory { @Override public TbRuleEngineProcessingDecision analyze(TbRuleEngineProcessingResult result) { if (result.isSuccess()) { - log.debug("[{}] The result of the msg pack processing is successful, going to proceed with processing of the following msgs", queueName); + log.trace("[{}] The result of the msg pack processing is successful, going to proceed with processing of the following msgs", queueName); return new TbRuleEngineProcessingDecision(true, null); } else { if (retryCount == 0) { @@ -107,8 +107,8 @@ public class TbRuleEngineProcessingStrategyFactory { } if (retrySuccessful) { result.getSuccessMap().forEach(toReprocess::put); - } else if (log.isDebugEnabled() && !result.getSuccessMap().isEmpty()) { - log.debug("[{}] Skipped {} successful messages due to the processing strategy configuration", queueName, result.getSuccessMap().size()); + } else if (log.isTraceEnabled() && !result.getSuccessMap().isEmpty()) { + log.trace("[{}] Skipped {} successful messages due to the processing strategy configuration", queueName, result.getSuccessMap().size()); } if (CollectionUtils.isEmpty(toReprocess)) { if (log.isDebugEnabled()) {