From 9d25db32c68a58dc331524989750af002419104a Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 24 Mar 2025 15:33:03 +0100 Subject: [PATCH] Refactored initialization chain, PartitionChangeEvent should be processed after actors startup --- ...faultTbCalculatedFieldConsumerService.java | 18 ++-- .../DefaultTbRuleEngineConsumerService.java | 17 ++-- .../AbstractConsumerPartitionedService.java | 94 +++++++++++++++++++ 3 files changed, 114 insertions(+), 15 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerPartitionedService.java 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 9ae06309c3..bce1992932 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 @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.queue; -import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections4.CollectionUtils; @@ -55,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.AbstractConsumerService; +import org.thingsboard.server.service.queue.processing.AbstractConsumerPartitionedService; import org.thingsboard.server.service.queue.processing.IdMsgPair; import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService; @@ -72,7 +71,7 @@ import java.util.stream.Collectors; @Service @TbRuleEngineComponent @Slf4j -public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerService implements TbCalculatedFieldConsumerService { +public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerPartitionedService implements TbCalculatedFieldConsumerService { @Value("${queue.calculated_fields.poll_interval:25}") private long pollInterval; @@ -99,10 +98,8 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer this.stateService = stateService; } - @PostConstruct - public void init() { - super.init("tb-cf"); - + @Override + protected void doAfterStartUp() { var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); PartitionedQueueConsumerManager> eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) @@ -129,7 +126,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer } @Override - protected void onTbApplicationEvent(PartitionChangeEvent event) { + protected void processPartitionChangeEvent(PartitionChangeEvent event) { try { event.getNewPartitions().forEach((queueKey, partitions) -> { if (queueKey.getQueueName().equals(DataConstants.CF_QUEUE_NAME)) { @@ -146,6 +143,11 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer } } + @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( 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 9cc743e510..3a9be3874a 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 @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.queue; -import jakarta.annotation.PostConstruct; import lombok.extern.slf4j.Slf4j; import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.event.EventListener; @@ -50,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.AbstractConsumerService; +import org.thingsboard.server.service.queue.processing.AbstractConsumerPartitionedService; import org.thingsboard.server.service.queue.ruleengine.TbRuleEngineConsumerContext; import org.thingsboard.server.service.queue.ruleengine.TbRuleEngineQueueConsumerManager; import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService; @@ -67,7 +66,7 @@ import java.util.stream.Collectors; @Service @TbRuleEngineComponent @Slf4j -public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService implements TbRuleEngineConsumerService { +public class DefaultTbRuleEngineConsumerService extends AbstractConsumerPartitionedService implements TbRuleEngineConsumerService { private final TbRuleEngineConsumerContext ctx; private final QueueService queueService; @@ -93,9 +92,8 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< this.queueService = queueService; } - @PostConstruct - public void init() { - super.init("tb-rule-engine"); + @Override + protected void doAfterStartUp() { List queues = queueService.findAllQueues(); for (Queue configuration : queues) { if (partitionService.isManagedByCurrentService(configuration.getTenantId())) { @@ -106,7 +104,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< } @Override - protected void onTbApplicationEvent(PartitionChangeEvent event) { + protected void processPartitionChangeEvent(PartitionChangeEvent event) { event.getNewPartitions().forEach((queueKey, partitions) -> { if (DataConstants.CF_QUEUE_NAME.equals(queueKey.getQueueName()) || DataConstants.CF_STATES_QUEUE_NAME.equals(queueKey.getQueueName())) { return; @@ -138,6 +136,11 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< }); } + @Override + protected String getPrefix() { + return "tb-rule-engine"; + } + @Override protected void stopConsumers() { super.stopConsumers(); 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/AbstractConsumerPartitionedService.java new file mode 100644 index 0000000000..94b5390be1 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerPartitionedService.java @@ -0,0 +1,94 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.queue.processing; + +import jakarta.annotation.PostConstruct; +import org.springframework.context.ApplicationEventPublisher; +import org.thingsboard.server.actors.ActorSystemContext; +import org.thingsboard.server.dao.tenant.TbTenantProfileCache; +import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; +import org.thingsboard.server.queue.util.AfterStartUp; +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.security.auth.jwt.settings.JwtSettingsService; + +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; + +public abstract class AbstractConsumerPartitionedService extends AbstractConsumerService { + + private final Lock startupLock; + private volatile boolean consumersInitialized; + private PartitionChangeEvent lastPartitionChangeEvent; + + public AbstractConsumerPartitionedService(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 + public void init() { + super.init(getPrefix()); + } + + @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) + public void afterStartUp() { + super.afterStartUp(); + doAfterStartUp(); + startupLock.lock(); + try { + processPartitionChangeEvent(lastPartitionChangeEvent); + consumersInitialized = true; + } finally { + startupLock.unlock(); + } + } + + @Override + protected void onTbApplicationEvent(PartitionChangeEvent event) { + if (!consumersInitialized) { + startupLock.lock(); + try { + if (!consumersInitialized) { + lastPartitionChangeEvent = event; + return; + } + } finally { + startupLock.unlock(); + } + } + processPartitionChangeEvent(event); + } + + protected abstract void doAfterStartUp(); + + protected abstract void processPartitionChangeEvent(PartitionChangeEvent event); + + protected abstract String getPrefix(); + +}