diff --git a/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java b/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java index 44ca548a68..cd46c4ba96 100644 --- a/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java +++ b/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java @@ -147,9 +147,12 @@ public class DefaultEdqsService implements EdqsService { syncLock.lock(); try { EdqsSyncState syncState = getSyncState(); - if (syncState != null && syncState.getStatus() == EdqsSyncStatus.FINISHED) { - log.info("EDQS sync is already finished"); - return; + if (syncState != null) { + EdqsSyncStatus status = syncState.getStatus(); + if (status == EdqsSyncStatus.FINISHED || status == EdqsSyncStatus.FAILED) { + log.info("EDQS sync is already " + status + ", ignoring the msg"); + return; + } } saveSyncState(EdqsSyncStatus.STARTED); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java index bce1992932..125a06299d 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java @@ -54,7 +54,7 @@ import org.thingsboard.server.service.cf.CalculatedFieldCache; import org.thingsboard.server.service.cf.CalculatedFieldStateService; import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache; -import org.thingsboard.server.service.queue.processing.AbstractConsumerPartitionedService; +import org.thingsboard.server.service.queue.processing.AbstractPartitionBasedConsumerService; import org.thingsboard.server.service.queue.processing.IdMsgPair; import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService; @@ -71,7 +71,7 @@ import java.util.stream.Collectors; @Service @TbRuleEngineComponent @Slf4j -public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerPartitionedService implements TbCalculatedFieldConsumerService { +public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBasedConsumerService implements TbCalculatedFieldConsumerService { @Value("${queue.calculated_fields.poll_interval:25}") private long pollInterval; @@ -99,7 +99,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerPar } @Override - protected void doAfterStartUp() { + protected void onStartUp() { var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); PartitionedQueueConsumerManager> eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) @@ -126,7 +126,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerPar } @Override - protected void processPartitionChangeEvent(PartitionChangeEvent event) { + protected void onPartitionChangeEvent(PartitionChangeEvent event) { try { event.getNewPartitions().forEach((queueKey, partitions) -> { if (queueKey.getQueueName().equals(DataConstants.CF_QUEUE_NAME)) { @@ -143,11 +143,6 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerPar } } - @Override - protected String getPrefix() { - return "tb-cf"; - } - private void processMsgs(List> msgs, TbQueueConsumer> consumer, QueueConfig config) throws Exception { List> orderedMsgList = msgs.stream().map(msg -> new IdMsgPair<>(UUID.randomUUID(), msg)).toList(); ConcurrentMap> pendingMap = orderedMsgList.stream().collect( @@ -195,6 +190,11 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerPar return ServiceType.TB_RULE_ENGINE; } + @Override + protected String getPrefix() { + return "tb-cf"; + } + @Override protected long getNotificationPollDuration() { return pollInterval; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java index 3a9be3874a..98f857750a 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java @@ -49,7 +49,7 @@ import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.cf.CalculatedFieldCache; import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache; -import org.thingsboard.server.service.queue.processing.AbstractConsumerPartitionedService; +import org.thingsboard.server.service.queue.processing.AbstractPartitionBasedConsumerService; import org.thingsboard.server.service.queue.ruleengine.TbRuleEngineConsumerContext; import org.thingsboard.server.service.queue.ruleengine.TbRuleEngineQueueConsumerManager; import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService; @@ -66,7 +66,7 @@ import java.util.stream.Collectors; @Service @TbRuleEngineComponent @Slf4j -public class DefaultTbRuleEngineConsumerService extends AbstractConsumerPartitionedService implements TbRuleEngineConsumerService { +public class DefaultTbRuleEngineConsumerService extends AbstractPartitionBasedConsumerService implements TbRuleEngineConsumerService { private final TbRuleEngineConsumerContext ctx; private final QueueService queueService; @@ -93,7 +93,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerPartitio } @Override - protected void doAfterStartUp() { + protected void onStartUp() { List queues = queueService.findAllQueues(); for (Queue configuration : queues) { if (partitionService.isManagedByCurrentService(configuration.getTenantId())) { @@ -104,7 +104,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerPartitio } @Override - protected void processPartitionChangeEvent(PartitionChangeEvent event) { + protected void onPartitionChangeEvent(PartitionChangeEvent event) { event.getNewPartitions().forEach((queueKey, partitions) -> { if (DataConstants.CF_QUEUE_NAME.equals(queueKey.getQueueName()) || DataConstants.CF_STATES_QUEUE_NAME.equals(queueKey.getQueueName())) { return; @@ -136,11 +136,6 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerPartitio }); } - @Override - protected String getPrefix() { - return "tb-rule-engine"; - } - @Override protected void stopConsumers() { super.stopConsumers(); @@ -153,6 +148,11 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerPartitio return ServiceType.TB_RULE_ENGINE; } + @Override + protected String getPrefix() { + return "tb-rule-engine"; + } + @Override protected long getNotificationPollDuration() { return ctx.getPollDuration(); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerPartitionedService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractPartitionBasedConsumerService.java similarity index 62% rename from application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerPartitionedService.java rename to application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractPartitionBasedConsumerService.java index 94b5390be1..97aa81d41c 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerPartitionedService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractPartitionBasedConsumerService.java @@ -31,24 +31,22 @@ import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsServ import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; -public abstract class AbstractConsumerPartitionedService extends AbstractConsumerService { +public abstract class AbstractPartitionBasedConsumerService extends AbstractConsumerService { - private final Lock startupLock; - private volatile boolean consumersInitialized; + private final Lock startupLock = new ReentrantLock(); + private volatile boolean started = false; private PartitionChangeEvent lastPartitionChangeEvent; - public AbstractConsumerPartitionedService(ActorSystemContext actorContext, - TbTenantProfileCache tenantProfileCache, - TbDeviceProfileCache deviceProfileCache, - TbAssetProfileCache assetProfileCache, - CalculatedFieldCache calculatedFieldCache, - TbApiUsageStateService apiUsageStateService, - PartitionService partitionService, - ApplicationEventPublisher eventPublisher, - JwtSettingsService jwtSettingsService) { + public AbstractPartitionBasedConsumerService(ActorSystemContext actorContext, + TbTenantProfileCache tenantProfileCache, + TbDeviceProfileCache deviceProfileCache, + TbAssetProfileCache assetProfileCache, + CalculatedFieldCache calculatedFieldCache, + TbApiUsageStateService apiUsageStateService, + PartitionService partitionService, + ApplicationEventPublisher eventPublisher, + JwtSettingsService jwtSettingsService) { super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService, eventPublisher, jwtSettingsService); - this.startupLock = new ReentrantLock(); - this.consumersInitialized = false; } @PostConstruct @@ -57,13 +55,14 @@ public abstract class AbstractConsumerPartitionedService