diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java index 825b68dec2..19d2b74bb5 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java @@ -91,7 +91,7 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta } } }) - .consumerCreator((queueConfig, partitionId) -> queueFactory.createCalculatedFieldStateConsumer()) + .consumerCreator((queueConfig, tpi) -> queueFactory.createCalculatedFieldStateConsumer()) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(eventConsumer.getConsumerExecutor()) .scheduler(eventConsumer.getScheduler()) 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 6adae15fe6..f108949efa 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,6 +17,7 @@ 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; @@ -35,7 +36,6 @@ import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; @@ -43,7 +43,6 @@ 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; @@ -61,7 +60,6 @@ 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; @@ -85,8 +83,6 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa private final CalculatedFieldStateService stateService; private final CalculatedFieldEntityProfileCache entityProfileCache; - private final ConcurrentMap>> consumers = new ConcurrentHashMap<>(); - public DefaultTbCalculatedFieldConsumerService(TbRuleEngineQueueFactory tbQueueFactory, ActorSystemContext actorContext, TbDeviceProfileCache deviceProfileCache, @@ -109,32 +105,18 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa @Override protected void onStartUp() { var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); - createConsumer(queueKey); - } - - private PartitionedQueueConsumerManager> createConsumer(QueueKey queueKey) { - String topic = partitionService.getTopic(queueKey); var eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) - .topic(topic) + .topic(partitionService.getTopic(queueKey)) .pollInterval(pollInterval) .msgPackProcessor(this::processMsgs) - .consumerCreator((queueConfig, partitionId) -> { - TopicPartitionInfo tpi = TopicPartitionInfo.builder() - .tenantId(queueKey.getTenantId()) - .topic(partitionService.getTopic(queueKey)) - .partition(partitionId) - .build(); - return queueFactory.createToCalculatedFieldMsgConsumer(tpi, partitionId); - }) + .consumerCreator((queueConfig, tpi) -> queueFactory.createToCalculatedFieldMsgConsumer(tpi)) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) .taskExecutor(mgmtExecutor) .build(); stateService.init(eventConsumer); - consumers.put(queueKey, eventConsumer); - return eventConsumer; } @PreDestroy @@ -151,26 +133,11 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa protected void onPartitionChangeEvent(PartitionChangeEvent event) { try { 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())) { - var consumer = Optional.ofNullable(consumers.get(queueKey)).orElseGet(() -> createConsumer(queueKey)); - if (consumer != null) { - stateService.restore(queueKey, partitions); - } - } + if (DataConstants.CF_QUEUE_NAME.equals(queueKey.getQueueName())) { + 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()); @@ -265,19 +232,13 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa if (event.getEvent() == ComponentLifecycleEvent.DELETED) { entityProfileCache.removeTenant(event.getTenantId()); - 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); - } - }); - }); + 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())); } } else if (event.getEntityId().getEntityType() == EntityType.ASSET_PROFILE) { if (event.getEvent() == ComponentLifecycleEvent.DELETED) { diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index a3003ba6ff..17decae3a0 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -206,7 +206,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService queueFactory.createToCoreMsgConsumer()) + .consumerCreator((config, tpi) -> queueFactory.createToCoreMsgConsumer()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) .taskExecutor(mgmtExecutor) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java index fdaa2103e2..7415b0bdad 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java @@ -19,7 +19,6 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; -import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; import org.jetbrains.annotations.NotNull; @@ -45,12 +44,12 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg; 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.discovery.QueueKey; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.EdgeContextComponent; -import org.thingsboard.server.queue.common.consumer.MainQueueConsumerManager; import org.thingsboard.server.service.queue.processing.AbstractConsumerService; import org.thingsboard.server.service.queue.processing.IdMsgPair; @@ -104,7 +103,7 @@ public class DefaultTbEdgeConsumerService extends AbstractConsumerService queueFactory.createEdgeMsgConsumer()) + .consumerCreator((config, tpi) -> queueFactory.createEdgeMsgConsumer()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) .taskExecutor(mgmtExecutor) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java index c7f7f600a7..c147836dab 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java @@ -72,7 +72,12 @@ public class TbRuleEngineQueueConsumerManager extends MainQueueConsumerManager { + Integer partitionId = topicPartitionInfo.getPartition().orElse(-1); + return ctx.getQueueFactory().createToRuleEngineMsgConsumer(queueConfig, partitionId); + }, + consumerExecutor, scheduler, taskExecutor, null); this.ctx = ctx; this.stats = new TbRuleEngineConsumerStats(queueKey, ctx.getStatsFactory()); } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java index 07575220eb..78d17ef368 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java @@ -139,7 +139,7 @@ public class EdqsProcessor implements TbQueueHandler, } consumer.commit(); }) - .consumerCreator((config, partitionId) -> queueFactory.createEdqsEventsConsumer()) + .consumerCreator((config, tpi) -> queueFactory.createEdqsEventsConsumer()) .queueAdmin(queueFactory.getEdqsQueueAdmin()) .consumerExecutor(consumersExecutor) .taskExecutor(taskExecutor) diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java index efdb1ead1c..0efe6e7d3b 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java @@ -59,7 +59,8 @@ public class KafkaEdqsStateService implements EdqsStateService { private final EdqsConfig config; private final EdqsPartitionService partitionService; private final KafkaEdqsQueueFactory queueFactory; - @Autowired @Lazy + @Autowired + @Lazy private EdqsProcessor edqsProcessor; private PartitionedQueueConsumerManager> stateConsumer; @@ -93,7 +94,7 @@ public class KafkaEdqsStateService implements EdqsStateService { } consumer.commit(); }) - .consumerCreator((config, partitionId) -> queueFactory.createEdqsStateConsumer()) + .consumerCreator((config, tpi) -> queueFactory.createEdqsStateConsumer()) .queueAdmin(queueAdmin) .consumerExecutor(eventConsumer.getConsumerExecutor()) .taskExecutor(eventConsumer.getTaskExecutor()) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java index ef9728344c..2233855a37 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java @@ -54,7 +54,7 @@ public class MainQueueConsumerManager msgPackProcessor; - protected final BiFunction> consumerCreator; + protected final BiFunction> consumerCreator; @Getter protected final ExecutorService consumerExecutor; @Getter @@ -74,7 +74,7 @@ public class MainQueueConsumerManager msgPackProcessor, - BiFunction> consumerCreator, + BiFunction> consumerCreator, ExecutorService consumerExecutor, ScheduledExecutorService scheduler, ExecutorService taskExecutor, @@ -313,7 +313,7 @@ public class MainQueueConsumerManager onStop.accept(tpi) : null; TbQueueConsumerTask consumer = new TbQueueConsumerTask<>(key, () -> { - TbQueueConsumer queueConsumer = consumerCreator.apply(config, partitionId); + TbQueueConsumer queueConsumer = consumerCreator.apply(config, tpi); if (startOffsetProvider != null && queueConsumer instanceof TbKafkaConsumerTemplate kafkaConsumer) { kafkaConsumer.setStartOffsetProvider(startOffsetProvider); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java index 0de1e53753..1b19fcab0e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java @@ -45,7 +45,7 @@ public class PartitionedQueueConsumerManager extends MainQ @Builder(builderMethodName = "create") // not to conflict with super.builder() public PartitionedQueueConsumerManager(QueueKey queueKey, String topic, long pollInterval, MsgPackProcessor msgPackProcessor, - BiFunction> consumerCreator, TbQueueAdmin queueAdmin, + BiFunction> consumerCreator, TbQueueAdmin queueAdmin, ExecutorService consumerExecutor, ScheduledExecutorService scheduler, ExecutorService taskExecutor, Consumer uncaughtErrorHandler) { super(queueKey, QueueConfig.of(true, pollInterval), msgPackProcessor, consumerCreator, consumerExecutor, scheduler, taskExecutor, uncaughtErrorHandler); 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 9fbd75d347..085d04f28c 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 @@ -134,7 +134,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi, Integer partitionId) { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi) { 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 198e2f7a83..1c46a624c3 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 @@ -105,7 +105,7 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi private final TbQueueAdmin housekeeperReprocessingAdmin; private final TbQueueAdmin edgeAdmin; private final TbQueueAdmin edgeEventAdmin; - private final TbKafkaAdmin cfAdmin; + private final TbQueueAdmin cfAdmin; private final TbQueueAdmin cfStateAdmin; private final TbQueueAdmin edqsEventsAdmin; private final TbKafkaAdmin edqsRequestsAdmin; @@ -515,9 +515,10 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi, Integer partitionId) { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi) { String queueName = DataConstants.CF_QUEUE_NAME; TenantId tenantId = tpi.getTenantId().orElse(TenantId.SYS_TENANT_ID); + Integer partitionId = tpi.getPartition().orElseThrow(() -> new IllegalArgumentException("PartitionId is required.")); 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 f387d1af11..112af9f646 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 @@ -93,7 +93,7 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { private final TbQueueAdmin housekeeperAdmin; private final TbQueueAdmin edgeAdmin; private final TbQueueAdmin edgeEventAdmin; - private final TbKafkaAdmin cfAdmin; + private final TbQueueAdmin cfAdmin; private final TbQueueAdmin cfStateAdmin; private final TbQueueAdmin edqsEventsAdmin; private final AtomicLong consumerCount = new AtomicLong(); @@ -316,9 +316,10 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi, Integer partitionId) { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi) { String queueName = DataConstants.CF_QUEUE_NAME; TenantId tenantId = tpi.getTenantId().orElse(TenantId.SYS_TENANT_ID); + Integer partitionId = tpi.getPartition().orElseThrow(() -> new IllegalArgumentException("PartitionId is required.")); 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 3980aaa632..18bb6db14a 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 @@ -122,7 +122,7 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate(); - TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi, Integer partitionId); + TbQueueConsumer> createToCalculatedFieldMsgConsumer(TopicPartitionInfo tpi); TbQueueAdmin getCalculatedFieldQueueAdmin();