From 434a534ba35005bc7c695ef26049cc2c86b6cdd9 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Tue, 8 Apr 2025 15:05:19 +0300 Subject: [PATCH] created cf consumer using topic partition info --- ...faultTbCalculatedFieldConsumerService.java | 63 +++++++++++++------ .../InMemoryMonolithQueueFactory.java | 3 +- .../provider/KafkaMonolithQueueFactory.java | 8 ++- .../KafkaTbRuleEngineQueueFactory.java | 8 ++- .../provider/TbRuleEngineQueueFactory.java | 3 +- 5 files changed, 59 insertions(+), 26 deletions(-) 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 ef4c4485e6..6adae15fe6 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 @@ -17,7 +17,6 @@ package org.thingsboard.server.service.queue; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.collections4.CollectionUtils; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.event.EventListener; @@ -30,7 +29,6 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; -import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.data.queue.QueueConfig; import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; @@ -45,6 +43,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.common.consumer.MainQueueConsumerManager; import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.QueueKey; @@ -62,6 +61,7 @@ import org.thingsboard.server.service.queue.processing.IdMsgPair; import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService; import java.util.List; +import java.util.Optional; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -84,7 +84,8 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa private final TbRuleEngineQueueFactory queueFactory; private final CalculatedFieldStateService stateService; private final CalculatedFieldEntityProfileCache entityProfileCache; - private final QueueService queueService; + + private final ConcurrentMap>> consumers = new ConcurrentHashMap<>(); public DefaultTbCalculatedFieldConsumerService(TbRuleEngineQueueFactory tbQueueFactory, ActorSystemContext actorContext, @@ -97,35 +98,42 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa JwtSettingsService jwtSettingsService, CalculatedFieldCache calculatedFieldCache, CalculatedFieldStateService stateService, - CalculatedFieldEntityProfileCache entityProfileCache, - QueueService queueService) { + CalculatedFieldEntityProfileCache entityProfileCache) { super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService, eventPublisher, jwtSettingsService); this.queueFactory = tbQueueFactory; this.stateService = stateService; this.entityProfileCache = entityProfileCache; - this.queueService = queueService; } @Override protected void onStartUp() { var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); - Queue queue = queueService.findQueueByTenantIdAndName(queueKey.getTenantId(), queueKey.getQueueName()); - createConsumer(queue, queueKey); + createConsumer(queueKey); } - private PartitionedQueueConsumerManager> createConsumer(Queue queue, QueueKey queueKey) { + private PartitionedQueueConsumerManager> createConsumer(QueueKey queueKey) { + String topic = partitionService.getTopic(queueKey); var eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) + .topic(topic) .pollInterval(pollInterval) .msgPackProcessor(this::processMsgs) - .consumerCreator((queueConfig, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer(queue, partitionId)) + .consumerCreator((queueConfig, partitionId) -> { + TopicPartitionInfo tpi = TopicPartitionInfo.builder() + .tenantId(queueKey.getTenantId()) + .topic(partitionService.getTopic(queueKey)) + .partition(partitionId) + .build(); + return queueFactory.createToCalculatedFieldMsgConsumer(tpi, partitionId); + }) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) .taskExecutor(mgmtExecutor) .build(); stateService.init(eventConsumer); + consumers.put(queueKey, eventConsumer); return eventConsumer; } @@ -145,11 +153,24 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa event.getNewPartitions().forEach((queueKey, partitions) -> { if (DataConstants.CF_QUEUE_NAME.equals(queueKey.getQueueName()) || DataConstants.CF_STATES_QUEUE_NAME.equals(queueKey.getQueueName())) { if (partitionService.isManagedByCurrentService(queueKey.getTenantId())) { - stateService.restore(queueKey, partitions); + var consumer = Optional.ofNullable(consumers.get(queueKey)).orElseGet(() -> createConsumer(queueKey)); + if (consumer != null) { + stateService.restore(queueKey, partitions); + } } } }); + consumers.keySet().stream() + .collect(Collectors.groupingBy(QueueKey::getTenantId)) + .forEach((tenantId, queueKeys) -> { + if (!partitionService.isManagedByCurrentService(tenantId)) { + queueKeys.forEach(queueKey -> { + Optional.ofNullable(consumers.remove(queueKey)).ifPresent(MainQueueConsumerManager::stop); + }); + } + }); + // Cleanup old entities after corresponding consumers are stopped. // Any periodic tasks need to check that the entity is still managed by the current server before processing. actorContext.tell(new CalculatedFieldPartitionChangeMsg()); @@ -244,13 +265,19 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa if (event.getEvent() == ComponentLifecycleEvent.DELETED) { entityProfileCache.removeTenant(event.getTenantId()); - Set partitions = stateService.getPartitions(); - if (CollectionUtils.isEmpty(partitions)) { - return; - } - stateService.delete(partitions.stream() - .filter(tpi -> tpi.getTenantId().isPresent() && tpi.getTenantId().get().equals(event.getTenantId())) - .collect(Collectors.toSet())); + List toRemove = consumers.keySet().stream() + .filter(queueKey -> queueKey.getTenantId().equals(event.getTenantId())) + .toList(); + toRemove.forEach(queueKey -> { + Optional.ofNullable(consumers.remove(queueKey)).ifPresent(consumer -> { + Set partitions = stateService.getPartitions().stream() + .filter(tpi -> tpi.getTenantId().isPresent() && tpi.getTenantId().get().equals(event.getTenantId())) + .collect(Collectors.toSet()); + if (!partitions.isEmpty()) { + consumer.delete(partitions); + } + }); + }); } } else if (event.getEntityId().getEntityType() == EntityType.ASSET_PROFILE) { if (event.getEvent() == ComponentLifecycleEvent.DELETED) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java index 2a8403040d..9fbd75d347 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java @@ -22,6 +22,7 @@ import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; @@ -133,7 +134,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi, Integer partitionId) { return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java index 9a52c775cb..198e2f7a83 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java @@ -20,10 +20,12 @@ import jakarta.annotation.PreDestroy; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.context.annotation.Bean; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; @@ -513,9 +515,9 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) { - String queueName = queue.getName(); - TenantId tenantId = queue.getTenantId(); + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi, Integer partitionId) { + String queueName = DataConstants.CF_QUEUE_NAME; + TenantId tenantId = tpi.getTenantId().orElse(TenantId.SYS_TENANT_ID); String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java index d1e3528189..f387d1af11 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java @@ -20,9 +20,11 @@ import jakarta.annotation.PreDestroy; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.context.annotation.Bean; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; @@ -314,9 +316,9 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) { - String queueName = queue.getName(); - TenantId tenantId = queue.getTenantId(); + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi, Integer partitionId) { + String queueName = DataConstants.CF_QUEUE_NAME; + TenantId tenantId = tpi.getTenantId().orElse(TenantId.SYS_TENANT_ID); String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java index 5cabfe9cbb..3980aaa632 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java @@ -16,6 +16,7 @@ package org.thingsboard.server.queue.provider; import org.thingsboard.server.common.data.queue.Queue; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; @@ -121,7 +122,7 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate(); - TbQueueConsumer> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId); + TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi, Integer partitionId); TbQueueAdmin getCalculatedFieldQueueAdmin();