diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java index a8b0ccf30e..7cb9e718f6 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.cf; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.exception.CalculatedFieldStateException; @@ -29,7 +30,7 @@ import java.util.Set; public interface CalculatedFieldStateService { - void init(PartitionedQueueConsumerManager> eventConsumer); + void init(TenantId tenantId, PartitionedQueueConsumerManager> eventConsumer); void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) throws CalculatedFieldStateException; 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 37f79edf6e..1eb8d0acaa 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 @@ -67,7 +67,7 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta private final AtomicInteger counter = new AtomicInteger(); @Override - public void init(PartitionedQueueConsumerManager> eventConsumer) { + public void init(TenantId tenantId, PartitionedQueueConsumerManager> eventConsumer) { var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); PartitionedQueueConsumerManager> stateConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) @@ -91,7 +91,7 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta } } }) - .consumerCreator((config, partitionId) -> queueFactory.createCalculatedFieldStateConsumer()) + .consumerCreator((config, partitionId) -> queueFactory.createCalculatedFieldStateConsumer(tenantId)) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(eventConsumer.getConsumerExecutor()) .scheduler(eventConsumer.getScheduler()) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java index 3d2f522440..fc30822b85 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java @@ -21,6 +21,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -45,8 +46,8 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS private final CfRocksDb cfRocksDb; @Override - public void init(PartitionedQueueConsumerManager> eventConsumer) { - super.stateServices.put(new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME), new DefaultQueueStateService<>(eventConsumer)); + public void init(TenantId tenantId, PartitionedQueueConsumerManager> eventConsumer) { + super.stateServices.put(new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId), new DefaultQueueStateService<>(eventConsumer)); } @Override 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 d4a88c2fe6..9e0376b483 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 @@ -45,6 +45,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 +63,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; @@ -88,7 +90,6 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa private final ConcurrentMap>> consumers = new ConcurrentHashMap<>(); - public DefaultTbCalculatedFieldConsumerService(TbRuleEngineQueueFactory tbQueueFactory, ActorSystemContext actorContext, TbDeviceProfileCache deviceProfileCache, @@ -115,19 +116,19 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa PageDataIterable iterator = new PageDataIterable<>(tenantService::findTenantsIds, 1024); for (TenantId tenantId : iterator) { if (partitionService.isManagedByCurrentService(tenantId)) { - stateService.init(createConsumer(tenantId)); + stateService.init(tenantId, createConsumer(tenantId)); } } } private PartitionedQueueConsumerManager> createConsumer(TenantId tenantId) { - QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); + QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId); var eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) - .topic(partitionService.getTopic(queueKey)) + .topic(partitionService.getTopic(new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME))) .pollInterval(pollInterval) .msgPackProcessor(this::processMsgs) - .consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer(tenantId, partitionId)) + .consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer(tenantId)) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) @@ -151,11 +152,30 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa protected void onPartitionChangeEvent(PartitionChangeEvent event) { try { event.getNewPartitions().forEach((queueKey, partitions) -> { - if (queueKey.getQueueName().equals(DataConstants.CF_QUEUE_NAME)) { - stateService.restore(queueKey, partitions); + if (!queueKey.getQueueName().equals(DataConstants.CF_QUEUE_NAME)) { + return; + } + if (partitionService.isManagedByCurrentService(queueKey.getTenantId())) { + var consumer = Optional.ofNullable(consumers.get(queueKey)).orElseGet(() -> { + var newConsumer = createConsumer(queueKey.getTenantId()); + stateService.init(queueKey.getTenantId(), newConsumer); + return newConsumer; + }); + if (consumer != null) { + stateService.restore(queueKey, partitions); + // eventConsumer's partitions will be updated by stateService + } } }); - // eventConsumer's partitions will be updated by stateService + 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. @@ -250,23 +270,23 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa if (event.getEntityId().getEntityType() == EntityType.TENANT) { if (event.getEvent() == ComponentLifecycleEvent.DELETED) { entityProfileCache.removeTenant(event.getTenantId()); - consumers.keySet().removeIf(queueKey -> { - boolean toRemove = queueKey.getTenantId().equals(event.getTenantId()); - if (toRemove) { + + 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(queueKey); if (!CollectionUtils.isEmpty(partitions)) { stateService.delete(queueKey, partitions); } - var consumer = consumers.get(queueKey); - if (consumer != null) { - consumer.stop(); - } - } - return toRemove; + consumer.stop(); + }); }); } else if (event.getEvent() == ComponentLifecycleEvent.CREATED) { if (partitionService.isManagedByCurrentService(event.getTenantId())) { - stateService.init(createConsumer(event.getTenantId())); + stateService.init(event.getTenantId(), createConsumer(event.getTenantId())); } } } else if (event.getEntityId().getEntityType() == EntityType.ASSET_PROFILE) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java index 32537edba5..7c2d0b0240 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java @@ -18,6 +18,7 @@ package org.thingsboard.server.queue.discovery.event; import lombok.Getter; import lombok.ToString; import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.queue.discovery.QueueKey; 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 94d53e4864..e825a1424d 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(TenantId tenantId, Integer partitionId) { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId) { return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); } @@ -154,7 +154,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE } @Override - public TbQueueConsumer> createCalculatedFieldStateConsumer() { + public TbQueueConsumer> createCalculatedFieldStateConsumer(TenantId tenantId) { return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getStateTopic())); } 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 4cb8fc5793..ca41ee221d 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 @@ -514,12 +514,12 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId, Integer partitionId) { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId) { String queueName = DataConstants.CF_QUEUE_NAME; - String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId); + String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId - groupId, partitionId); + groupId, null); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); @@ -572,18 +572,25 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - public TbQueueConsumer> createCalculatedFieldStateConsumer() { - return TbKafkaConsumerTemplate.>builder() - .settings(kafkaSettings) - .topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic())) - .readFromBeginning(true) - .stopWhenRead(true) - .clientId("monolith-calculated-field-state-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()) - .groupId(topicService.buildTopicName("monolith-calculated-field-state-consumer")) - .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), msg.getData() != null ? CalculatedFieldStateProto.parseFrom(msg.getData()) : null, msg.getHeaders())) - .admin(cfStateAdmin) - .statsService(consumerStatsService) - .build(); + public TbQueueConsumer> createCalculatedFieldStateConsumer(TenantId tenantId) { + String queueName = DataConstants.CF_STATES_QUEUE_NAME; + String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); + + cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId + groupId, null); + + TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); + consumerBuilder.settings(kafkaSettings); + consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic())); + consumerBuilder.readFromBeginning(true); + consumerBuilder.stopWhenRead(true); + consumerBuilder.clientId("cf-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); + consumerBuilder.groupId(groupId); + consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), msg.getData() != null ? CalculatedFieldStateProto.parseFrom(msg.getData()) : null, msg.getHeaders())); + consumerBuilder.admin(cfStateAdmin); + consumerBuilder.statsService(consumerStatsService); + + return consumerBuilder.build(); } @Override 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 ae565a0827..84ad126ff2 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 @@ -315,17 +315,17 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId, Integer partitionId) { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId) { String queueName = DataConstants.CF_QUEUE_NAME; - String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId); + String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId - groupId, partitionId); + groupId, null); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); - consumerBuilder.clientId("cf-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); + consumerBuilder.clientId("cf-" + queueName + "-consumer-" + tenantId + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); consumerBuilder.groupId(groupId); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCalculatedFieldMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(cfAdmin); @@ -372,18 +372,24 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - public TbQueueConsumer> createCalculatedFieldStateConsumer() { - return TbKafkaConsumerTemplate.>builder() - .settings(kafkaSettings) - .topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic())) - .readFromBeginning(true) - .stopWhenRead(true) - .clientId("tb-rule-engine-calculated-field-state-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()) - .groupId(topicService.buildTopicName("tb-rule-engine-calculated-field-state-consumer")) - .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), msg.getData() != null ? CalculatedFieldStateProto.parseFrom(msg.getData()) : null, msg.getHeaders())) - .admin(cfStateAdmin) - .statsService(consumerStatsService) - .build(); + public TbQueueConsumer> createCalculatedFieldStateConsumer(TenantId tenantId) { + String queueName = DataConstants.CF_STATES_QUEUE_NAME; + String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); + + cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId + groupId, null); + + TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); + consumerBuilder.settings(kafkaSettings); + consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic())); + consumerBuilder.readFromBeginning(true); + consumerBuilder.stopWhenRead(true); + consumerBuilder.clientId("cf-" + queueName + "-consumer-" + tenantId + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); + consumerBuilder.groupId(groupId); + consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), msg.getData() != null ? CalculatedFieldStateProto.parseFrom(msg.getData()) : null, msg.getHeaders())); + consumerBuilder.admin(cfStateAdmin); + consumerBuilder.statsService(consumerStatsService); + return consumerBuilder.build(); } @Override 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 b71f6a0471..7076ceb0b0 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(TenantId tenantId, Integer partitionId); + TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId); TbQueueAdmin getCalculatedFieldQueueAdmin(); @@ -132,7 +132,7 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory TbQueueProducer> createToCalculatedFieldNotificationMsgProducer(); - TbQueueConsumer> createCalculatedFieldStateConsumer(); + TbQueueConsumer> createCalculatedFieldStateConsumer(TenantId tenantId); TbQueueProducer> createCalculatedFieldStateProducer();