diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java index 1347705a23..91c08ab6e0 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java @@ -30,9 +30,7 @@ import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; import java.util.Collection; -import java.util.Map; import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto; @@ -43,7 +41,7 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF @Autowired private ActorSystemContext actorSystemContext; - protected Map, TbProtoQueueMsg>> stateServices = new ConcurrentHashMap<>(); + protected QueueStateService, TbProtoQueueMsg> stateService; @Override public final void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) { @@ -74,22 +72,22 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF @Override public void restore(QueueKey queueKey, Set partitions) { - stateServices.get(queueKey).update(queueKey, partitions); + stateService.update(queueKey, partitions); } @Override - public void delete(QueueKey queueKey, Set partitions) { - stateServices.get(queueKey).delete(partitions); + public void delete(Set partitions) { + stateService.delete(partitions); } @Override - public Set getPartitions(QueueKey queueKey) { - return stateServices.get(queueKey).getPartitions().values().stream().flatMap(Collection::stream).collect(Collectors.toSet()); + public Set getPartitions() { + return stateService.getPartitions().values().stream().flatMap(Collection::stream).collect(Collectors.toSet()); } @Override - public void stop(QueueKey queueKey) { - stateServices.get(queueKey).stop(); + public void stop() { + stateService.stop(); } } 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 7cb9e718f6..d0b34f18e8 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,7 +15,6 @@ */ 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; @@ -30,7 +29,7 @@ import java.util.Set; public interface CalculatedFieldStateService { - void init(TenantId tenantId, PartitionedQueueConsumerManager> eventConsumer); + void init(PartitionedQueueConsumerManager> eventConsumer); void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) throws CalculatedFieldStateException; @@ -38,10 +37,10 @@ public interface CalculatedFieldStateService { void restore(QueueKey queueKey, Set partitions); - void delete(QueueKey queueKey, Set partitions); + void delete(Set partitions); - Set getPartitions(QueueKey queueKey); + Set getPartitions(); - void stop(QueueKey queueKey); + void stop(); } 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 1eb8d0acaa..825b68dec2 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,8 +67,8 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta private final AtomicInteger counter = new AtomicInteger(); @Override - public void init(TenantId tenantId, PartitionedQueueConsumerManager> eventConsumer) { - var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); + public void init(PartitionedQueueConsumerManager> eventConsumer) { + var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_STATES_QUEUE_NAME); PartitionedQueueConsumerManager> stateConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) .topic(partitionService.getTopic(queueKey)) @@ -91,16 +91,16 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta } } }) - .consumerCreator((config, partitionId) -> queueFactory.createCalculatedFieldStateConsumer(tenantId)) + .consumerCreator((queueConfig, partitionId) -> queueFactory.createCalculatedFieldStateConsumer()) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(eventConsumer.getConsumerExecutor()) .scheduler(eventConsumer.getScheduler()) .taskExecutor(eventConsumer.getTaskExecutor()) .build(); - super.stateServices.put(queueKey, KafkaQueueStateService., TbProtoQueueMsg>builder() + super.stateService = KafkaQueueStateService., TbProtoQueueMsg>builder() .eventConsumer(eventConsumer) .stateConsumer(stateConsumer) - .build()); + .build(); this.stateProducer = (TbKafkaProducerTemplate>) queueFactory.createCalculatedFieldStateProducer(); } @@ -148,8 +148,8 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta } @Override - public void stop(QueueKey queueKey) { - super.stop(queueKey); + public void stop() { + super.stop(); stateProducer.stop(); } 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 fc30822b85..9dc6139ca5 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 @@ -20,9 +20,6 @@ import lombok.RequiredArgsConstructor; 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; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; @@ -46,8 +43,8 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS private final CfRocksDb cfRocksDb; @Override - public void init(TenantId tenantId, PartitionedQueueConsumerManager> eventConsumer) { - super.stateServices.put(new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId), new DefaultQueueStateService<>(eventConsumer)); + public void init(PartitionedQueueConsumerManager> eventConsumer) { + super.stateService = new DefaultQueueStateService<>(eventConsumer); } @Override @@ -64,7 +61,7 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS @Override public void restore(QueueKey queueKey, Set partitions) { - if (stateServices.get(queueKey).getPartitions().isEmpty()) { + if (stateService.getPartitions().isEmpty()) { cfRocksDb.forEach((key, value) -> { try { processRestoredState(CalculatedFieldStateProto.parseFrom(value)); 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 9e0376b483..ef4c4485e6 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 @@ -29,23 +29,22 @@ import org.thingsboard.server.common.data.DataConstants; 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.page.PageDataIterable; 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; 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.dao.tenant.TenantService; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; 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; @@ -63,7 +62,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; @@ -84,11 +82,9 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa private long packProcessingTimeout; private final TbRuleEngineQueueFactory queueFactory; - private final TenantService tenantService; private final CalculatedFieldStateService stateService; private final CalculatedFieldEntityProfileCache entityProfileCache; - - private final ConcurrentMap>> consumers = new ConcurrentHashMap<>(); + private final QueueService queueService; public DefaultTbCalculatedFieldConsumerService(TbRuleEngineQueueFactory tbQueueFactory, ActorSystemContext actorContext, @@ -102,39 +98,34 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa CalculatedFieldCache calculatedFieldCache, CalculatedFieldStateService stateService, CalculatedFieldEntityProfileCache entityProfileCache, - TenantService tenantService) { + QueueService queueService) { super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService, eventPublisher, jwtSettingsService); this.queueFactory = tbQueueFactory; this.stateService = stateService; this.entityProfileCache = entityProfileCache; - this.tenantService = tenantService; + this.queueService = queueService; } @Override protected void onStartUp() { - PageDataIterable iterator = new PageDataIterable<>(tenantService::findTenantsIds, 1024); - for (TenantId tenantId : iterator) { - if (partitionService.isManagedByCurrentService(tenantId)) { - stateService.init(tenantId, createConsumer(tenantId)); - } - } + var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); + Queue queue = queueService.findQueueByTenantIdAndName(queueKey.getTenantId(), queueKey.getQueueName()); + createConsumer(queue, queueKey); } - private PartitionedQueueConsumerManager> createConsumer(TenantId tenantId) { - QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId); + private PartitionedQueueConsumerManager> createConsumer(Queue queue, QueueKey queueKey) { var eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(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)) + .consumerCreator((queueConfig, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer(queue, partitionId)) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .consumerExecutor(consumersExecutor) .scheduler(scheduler) .taskExecutor(mgmtExecutor) .build(); - consumers.put(queueKey, eventConsumer); + stateService.init(eventConsumer); return eventConsumer; } @@ -152,30 +143,12 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa protected void onPartitionChangeEvent(PartitionChangeEvent event) { try { event.getNewPartitions().forEach((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) { + 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); - // 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. @@ -271,23 +244,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(queueKey); - if (!CollectionUtils.isEmpty(partitions)) { - stateService.delete(queueKey, partitions); - } - consumer.stop(); - }); - }); - } else if (event.getEvent() == ComponentLifecycleEvent.CREATED) { - if (partitionService.isManagedByCurrentService(event.getTenantId())) { - stateService.init(event.getTenantId(), createConsumer(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())); } } else if (event.getEntityId().getEntityType() == EntityType.ASSET_PROFILE) { if (event.getEvent() == ComponentLifecycleEvent.DELETED) { @@ -320,7 +283,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa @Override protected void stopConsumers() { super.stopConsumers(); - consumers.keySet().forEach(stateService::stop); // eventConsumer will be stopped by stateService + stateService.stop(); // eventConsumer will be stopped by stateService } } 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 e825a1424d..2a8403040d 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 @@ -20,7 +20,6 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; -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.gen.js.JsInvokeProtos; @@ -134,7 +133,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId) { + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) { return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); } @@ -154,7 +153,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE } @Override - public TbQueueConsumer> createCalculatedFieldStateConsumer(TenantId tenantId) { + public TbQueueConsumer> createCalculatedFieldStateConsumer() { 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 ca41ee221d..9a52c775cb 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,7 +20,6 @@ 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; @@ -514,12 +513,10 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId) { - String queueName = DataConstants.CF_QUEUE_NAME; - String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); - - cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId - groupId, null); + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) { + String queueName = queue.getName(); + TenantId tenantId = queue.getTenantId(); + String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); @@ -572,25 +569,18 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi } @Override - 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(); + 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(); } @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 84ad126ff2..d1e3528189 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,7 +20,6 @@ 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; @@ -315,17 +314,15 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - public TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId) { - String queueName = DataConstants.CF_QUEUE_NAME; - String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); - - cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId - groupId, null); + public TbQueueConsumer> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) { + String queueName = queue.getName(); + TenantId tenantId = queue.getTenantId(); + String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); - consumerBuilder.clientId("cf-" + queueName + "-consumer-" + tenantId + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); + consumerBuilder.clientId("cf-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); consumerBuilder.groupId(groupId); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCalculatedFieldMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(cfAdmin); @@ -372,24 +369,18 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { } @Override - 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(); + 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(); } @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 7076ceb0b0..5cabfe9cbb 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 @@ -15,7 +15,6 @@ */ package org.thingsboard.server.queue.provider; -import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; @@ -96,8 +95,8 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory /** * Used to consume messages by TB Rule Engine Service * - * @return * @param configuration + * @return */ TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration); @@ -105,9 +104,9 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory * Used to consume messages by TB Rule Engine Service * Intended usage for consumer per partition strategy * - * @return TbQueueConsumer * @param configuration - * @param partitionId as a suffix for consumer name + * @param partitionId as a suffix for consumer name + * @return TbQueueConsumer */ default TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration, Integer partitionId) { return createToRuleEngineMsgConsumer(configuration); @@ -122,7 +121,7 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate(); - TbQueueConsumer> createToCalculatedFieldMsgConsumer(TenantId tenantId); + TbQueueConsumer> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId); TbQueueAdmin getCalculatedFieldQueueAdmin(); @@ -132,7 +131,7 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory TbQueueProducer> createToCalculatedFieldNotificationMsgProducer(); - TbQueueConsumer> createCalculatedFieldStateConsumer(TenantId tenantId); + TbQueueConsumer> createCalculatedFieldStateConsumer(); TbQueueProducer> createCalculatedFieldStateProducer();