Browse Source

wip consumer groups

pull/13040/head
IrynaMatveieva 1 year ago
parent
commit
88faa2fb75
  1. 18
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java
  2. 9
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java
  3. 14
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java
  4. 9
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java
  5. 77
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java
  6. 5
      common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryMonolithQueueFactory.java
  7. 42
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaMonolithQueueFactory.java
  8. 43
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java
  9. 11
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java

18
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 org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState;
import java.util.Collection; import java.util.Collection;
import java.util.Map;
import java.util.Set; import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto; import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto;
@ -43,7 +41,7 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF
@Autowired @Autowired
private ActorSystemContext actorSystemContext; private ActorSystemContext actorSystemContext;
protected Map<QueueKey, QueueStateService<TbProtoQueueMsg<ToCalculatedFieldMsg>, TbProtoQueueMsg<CalculatedFieldStateProto>>> stateServices = new ConcurrentHashMap<>(); protected QueueStateService<TbProtoQueueMsg<ToCalculatedFieldMsg>, TbProtoQueueMsg<CalculatedFieldStateProto>> stateService;
@Override @Override
public final void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) { public final void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) {
@ -74,22 +72,22 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF
@Override @Override
public void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions) { public void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions) {
stateServices.get(queueKey).update(queueKey, partitions); stateService.update(queueKey, partitions);
} }
@Override @Override
public void delete(QueueKey queueKey, Set<TopicPartitionInfo> partitions) { public void delete(Set<TopicPartitionInfo> partitions) {
stateServices.get(queueKey).delete(partitions); stateService.delete(partitions);
} }
@Override @Override
public Set<TopicPartitionInfo> getPartitions(QueueKey queueKey) { public Set<TopicPartitionInfo> getPartitions() {
return stateServices.get(queueKey).getPartitions().values().stream().flatMap(Collection::stream).collect(Collectors.toSet()); return stateService.getPartitions().values().stream().flatMap(Collection::stream).collect(Collectors.toSet());
} }
@Override @Override
public void stop(QueueKey queueKey) { public void stop() {
stateServices.get(queueKey).stop(); stateService.stop();
} }
} }

9
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java

@ -15,7 +15,6 @@
*/ */
package org.thingsboard.server.service.cf; 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.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.exception.CalculatedFieldStateException; import org.thingsboard.server.exception.CalculatedFieldStateException;
@ -30,7 +29,7 @@ import java.util.Set;
public interface CalculatedFieldStateService { public interface CalculatedFieldStateService {
void init(TenantId tenantId, PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>> eventConsumer); void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>> eventConsumer);
void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) throws CalculatedFieldStateException; void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) throws CalculatedFieldStateException;
@ -38,10 +37,10 @@ public interface CalculatedFieldStateService {
void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions); void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions);
void delete(QueueKey queueKey, Set<TopicPartitionInfo> partitions); void delete(Set<TopicPartitionInfo> partitions);
Set<TopicPartitionInfo> getPartitions(QueueKey queueKey); Set<TopicPartitionInfo> getPartitions();
void stop(QueueKey queueKey); void stop();
} }

14
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(); private final AtomicInteger counter = new AtomicInteger();
@Override @Override
public void init(TenantId tenantId, PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>> eventConsumer) { public void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>> eventConsumer) {
var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_STATES_QUEUE_NAME);
PartitionedQueueConsumerManager<TbProtoQueueMsg<CalculatedFieldStateProto>> stateConsumer = PartitionedQueueConsumerManager.<TbProtoQueueMsg<CalculatedFieldStateProto>>create() PartitionedQueueConsumerManager<TbProtoQueueMsg<CalculatedFieldStateProto>> stateConsumer = PartitionedQueueConsumerManager.<TbProtoQueueMsg<CalculatedFieldStateProto>>create()
.queueKey(queueKey) .queueKey(queueKey)
.topic(partitionService.getTopic(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()) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin())
.consumerExecutor(eventConsumer.getConsumerExecutor()) .consumerExecutor(eventConsumer.getConsumerExecutor())
.scheduler(eventConsumer.getScheduler()) .scheduler(eventConsumer.getScheduler())
.taskExecutor(eventConsumer.getTaskExecutor()) .taskExecutor(eventConsumer.getTaskExecutor())
.build(); .build();
super.stateServices.put(queueKey, KafkaQueueStateService.<TbProtoQueueMsg<ToCalculatedFieldMsg>, TbProtoQueueMsg<CalculatedFieldStateProto>>builder() super.stateService = KafkaQueueStateService.<TbProtoQueueMsg<ToCalculatedFieldMsg>, TbProtoQueueMsg<CalculatedFieldStateProto>>builder()
.eventConsumer(eventConsumer) .eventConsumer(eventConsumer)
.stateConsumer(stateConsumer) .stateConsumer(stateConsumer)
.build()); .build();
this.stateProducer = (TbKafkaProducerTemplate<TbProtoQueueMsg<CalculatedFieldStateProto>>) queueFactory.createCalculatedFieldStateProducer(); this.stateProducer = (TbKafkaProducerTemplate<TbProtoQueueMsg<CalculatedFieldStateProto>>) queueFactory.createCalculatedFieldStateProducer();
} }
@ -148,8 +148,8 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta
} }
@Override @Override
public void stop(QueueKey queueKey) { public void stop() {
super.stop(queueKey); super.stop();
stateProducer.stop(); stateProducer.stop();
} }

9
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 lombok.extern.slf4j.Slf4j;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Service; 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.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto;
@ -46,8 +43,8 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS
private final CfRocksDb cfRocksDb; private final CfRocksDb cfRocksDb;
@Override @Override
public void init(TenantId tenantId, PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>> eventConsumer) { public void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>> eventConsumer) {
super.stateServices.put(new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId), new DefaultQueueStateService<>(eventConsumer)); super.stateService = new DefaultQueueStateService<>(eventConsumer);
} }
@Override @Override
@ -64,7 +61,7 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS
@Override @Override
public void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions) { public void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions) {
if (stateServices.get(queueKey).getPartitions().isEmpty()) { if (stateService.getPartitions().isEmpty()) {
cfRocksDb.forEach((key, value) -> { cfRocksDb.forEach((key, value) -> {
try { try {
processRestoredState(CalculatedFieldStateProto.parseFrom(value)); processRestoredState(CalculatedFieldStateProto.parseFrom(value));

77
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.EntityType;
import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.TenantId; 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.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.data.queue.Queue;
import org.thingsboard.server.common.data.queue.QueueConfig; import org.thingsboard.server.common.data.queue.QueueConfig;
import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; 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.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; 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.common.consumer.PartitionedQueueConsumerManager;
import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.QueueKey; 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 org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService;
import java.util.List; import java.util.List;
import java.util.Optional;
import java.util.Set; import java.util.Set;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ -84,11 +82,9 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa
private long packProcessingTimeout; private long packProcessingTimeout;
private final TbRuleEngineQueueFactory queueFactory; private final TbRuleEngineQueueFactory queueFactory;
private final TenantService tenantService;
private final CalculatedFieldStateService stateService; private final CalculatedFieldStateService stateService;
private final CalculatedFieldEntityProfileCache entityProfileCache; private final CalculatedFieldEntityProfileCache entityProfileCache;
private final QueueService queueService;
private final ConcurrentMap<QueueKey, PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>>> consumers = new ConcurrentHashMap<>();
public DefaultTbCalculatedFieldConsumerService(TbRuleEngineQueueFactory tbQueueFactory, public DefaultTbCalculatedFieldConsumerService(TbRuleEngineQueueFactory tbQueueFactory,
ActorSystemContext actorContext, ActorSystemContext actorContext,
@ -102,39 +98,34 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa
CalculatedFieldCache calculatedFieldCache, CalculatedFieldCache calculatedFieldCache,
CalculatedFieldStateService stateService, CalculatedFieldStateService stateService,
CalculatedFieldEntityProfileCache entityProfileCache, CalculatedFieldEntityProfileCache entityProfileCache,
TenantService tenantService) { QueueService queueService) {
super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService, super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService,
eventPublisher, jwtSettingsService); eventPublisher, jwtSettingsService);
this.queueFactory = tbQueueFactory; this.queueFactory = tbQueueFactory;
this.stateService = stateService; this.stateService = stateService;
this.entityProfileCache = entityProfileCache; this.entityProfileCache = entityProfileCache;
this.tenantService = tenantService; this.queueService = queueService;
} }
@Override @Override
protected void onStartUp() { protected void onStartUp() {
PageDataIterable<TenantId> iterator = new PageDataIterable<>(tenantService::findTenantsIds, 1024); var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME);
for (TenantId tenantId : iterator) { Queue queue = queueService.findQueueByTenantIdAndName(queueKey.getTenantId(), queueKey.getQueueName());
if (partitionService.isManagedByCurrentService(tenantId)) { createConsumer(queue, queueKey);
stateService.init(tenantId, createConsumer(tenantId));
}
}
} }
private PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>> createConsumer(TenantId tenantId) { private PartitionedQueueConsumerManager<TbProtoQueueMsg<ToCalculatedFieldMsg>> createConsumer(Queue queue, QueueKey queueKey) {
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId);
var eventConsumer = PartitionedQueueConsumerManager.<TbProtoQueueMsg<ToCalculatedFieldMsg>>create() var eventConsumer = PartitionedQueueConsumerManager.<TbProtoQueueMsg<ToCalculatedFieldMsg>>create()
.queueKey(queueKey) .queueKey(queueKey)
.topic(partitionService.getTopic(new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME)))
.pollInterval(pollInterval) .pollInterval(pollInterval)
.msgPackProcessor(this::processMsgs) .msgPackProcessor(this::processMsgs)
.consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer(tenantId)) .consumerCreator((queueConfig, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer(queue, partitionId))
.queueAdmin(queueFactory.getCalculatedFieldQueueAdmin()) .queueAdmin(queueFactory.getCalculatedFieldQueueAdmin())
.consumerExecutor(consumersExecutor) .consumerExecutor(consumersExecutor)
.scheduler(scheduler) .scheduler(scheduler)
.taskExecutor(mgmtExecutor) .taskExecutor(mgmtExecutor)
.build(); .build();
consumers.put(queueKey, eventConsumer); stateService.init(eventConsumer);
return eventConsumer; return eventConsumer;
} }
@ -152,30 +143,12 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa
protected void onPartitionChangeEvent(PartitionChangeEvent event) { protected void onPartitionChangeEvent(PartitionChangeEvent event) {
try { try {
event.getNewPartitions().forEach((queueKey, partitions) -> { event.getNewPartitions().forEach((queueKey, partitions) -> {
if (!queueKey.getQueueName().equals(DataConstants.CF_QUEUE_NAME)) { if (DataConstants.CF_QUEUE_NAME.equals(queueKey.getQueueName()) || DataConstants.CF_STATES_QUEUE_NAME.equals(queueKey.getQueueName())) {
return; if (partitionService.isManagedByCurrentService(queueKey.getTenantId())) {
}
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); 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. // 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. // 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) { if (event.getEvent() == ComponentLifecycleEvent.DELETED) {
entityProfileCache.removeTenant(event.getTenantId()); entityProfileCache.removeTenant(event.getTenantId());
List<QueueKey> toRemove = consumers.keySet().stream() Set<TopicPartitionInfo> partitions = stateService.getPartitions();
.filter(queueKey -> queueKey.getTenantId().equals(event.getTenantId())) if (CollectionUtils.isEmpty(partitions)) {
.toList(); return;
toRemove.forEach(queueKey -> {
Optional.ofNullable(consumers.remove(queueKey)).ifPresent(consumer -> {
Set<TopicPartitionInfo> 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()));
} }
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) { } else if (event.getEntityId().getEntityType() == EntityType.ASSET_PROFILE) {
if (event.getEvent() == ComponentLifecycleEvent.DELETED) { if (event.getEvent() == ComponentLifecycleEvent.DELETED) {
@ -320,7 +283,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractPartitionBa
@Override @Override
protected void stopConsumers() { protected void stopConsumers() {
super.stopConsumers(); super.stopConsumers();
consumers.keySet().forEach(stateService::stop); // eventConsumer will be stopped by stateService stateService.stop(); // eventConsumer will be stopped by stateService
} }
} }

5
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.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.scheduling.annotation.Scheduled; import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component; 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.data.queue.Queue;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.js.JsInvokeProtos;
@ -134,7 +133,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<TransportProtos.ToCalculatedFieldMsg>> createToCalculatedFieldMsgConsumer(TenantId tenantId) { public TbQueueConsumer<TbProtoQueueMsg<TransportProtos.ToCalculatedFieldMsg>> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) {
return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getEventTopic()));
} }
@ -154,7 +153,7 @@ public class InMemoryMonolithQueueFactory implements TbCoreQueueFactory, TbRuleE
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateConsumer(TenantId tenantId) { public TbQueueConsumer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateConsumer() {
return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getStateTopic())); return new InMemoryTbQueueConsumer<>(storage, topicService.buildTopicName(calculatedFieldSettings.getStateTopic()));
} }

42
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.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
import org.springframework.stereotype.Component; 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.EdgeId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.data.queue.Queue;
@ -514,12 +513,10 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> createToCalculatedFieldMsgConsumer(TenantId tenantId) { public TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) {
String queueName = DataConstants.CF_QUEUE_NAME; String queueName = queue.getName();
String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); TenantId tenantId = queue.getTenantId();
String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId);
cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId
groupId, null);
TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder<TbProtoQueueMsg<ToCalculatedFieldMsg>> consumerBuilder = TbKafkaConsumerTemplate.builder(); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder<TbProtoQueueMsg<ToCalculatedFieldMsg>> consumerBuilder = TbKafkaConsumerTemplate.builder();
consumerBuilder.settings(kafkaSettings); consumerBuilder.settings(kafkaSettings);
@ -572,25 +569,18 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateConsumer(TenantId tenantId) { public TbQueueConsumer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateConsumer() {
String queueName = DataConstants.CF_STATES_QUEUE_NAME; return TbKafkaConsumerTemplate.<TbProtoQueueMsg<CalculatedFieldStateProto>>builder()
String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); .settings(kafkaSettings)
.topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic()))
cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId .readFromBeginning(true)
groupId, null); .stopWhenRead(true)
.clientId("monolith-calculated-field-state-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet())
TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder<TbProtoQueueMsg<CalculatedFieldStateProto>> consumerBuilder = TbKafkaConsumerTemplate.builder(); .groupId(topicService.buildTopicName("monolith-calculated-field-state-consumer"))
consumerBuilder.settings(kafkaSettings); .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), msg.getData() != null ? CalculatedFieldStateProto.parseFrom(msg.getData()) : null, msg.getHeaders()))
consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic())); .admin(cfStateAdmin)
consumerBuilder.readFromBeginning(true); .statsService(consumerStatsService)
consumerBuilder.stopWhenRead(true); .build();
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 @Override

43
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.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Bean;
import org.springframework.stereotype.Component; 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.id.TenantId;
import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.data.queue.Queue;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
@ -315,17 +314,15 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> createToCalculatedFieldMsgConsumer(TenantId tenantId) { public TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId) {
String queueName = DataConstants.CF_QUEUE_NAME; String queueName = queue.getName();
String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); TenantId tenantId = queue.getTenantId();
String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, partitionId);
cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId
groupId, null);
TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder<TbProtoQueueMsg<ToCalculatedFieldMsg>> consumerBuilder = TbKafkaConsumerTemplate.builder(); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder<TbProtoQueueMsg<ToCalculatedFieldMsg>> consumerBuilder = TbKafkaConsumerTemplate.builder();
consumerBuilder.settings(kafkaSettings); consumerBuilder.settings(kafkaSettings);
consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getEventTopic())); 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.groupId(groupId);
consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCalculatedFieldMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCalculatedFieldMsg.parseFrom(msg.getData()), msg.getHeaders()));
consumerBuilder.admin(cfAdmin); consumerBuilder.admin(cfAdmin);
@ -372,24 +369,18 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
} }
@Override @Override
public TbQueueConsumer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateConsumer(TenantId tenantId) { public TbQueueConsumer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateConsumer() {
String queueName = DataConstants.CF_STATES_QUEUE_NAME; return TbKafkaConsumerTemplate.<TbProtoQueueMsg<CalculatedFieldStateProto>>builder()
String groupId = topicService.buildConsumerGroupId("cf-", tenantId, queueName, null); .settings(kafkaSettings)
.topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic()))
cfAdmin.syncOffsets(topicService.buildConsumerGroupId("cf-", tenantId, queueName, null), // the fat groupId .readFromBeginning(true)
groupId, null); .stopWhenRead(true)
.clientId("tb-rule-engine-calculated-field-state-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet())
TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder<TbProtoQueueMsg<CalculatedFieldStateProto>> consumerBuilder = TbKafkaConsumerTemplate.builder(); .groupId(topicService.buildTopicName("tb-rule-engine-calculated-field-state-consumer"))
consumerBuilder.settings(kafkaSettings); .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), msg.getData() != null ? CalculatedFieldStateProto.parseFrom(msg.getData()) : null, msg.getHeaders()))
consumerBuilder.topic(topicService.buildTopicName(calculatedFieldSettings.getStateTopic())); .admin(cfStateAdmin)
consumerBuilder.readFromBeginning(true); .statsService(consumerStatsService)
consumerBuilder.stopWhenRead(true); .build();
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 @Override

11
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbRuleEngineQueueFactory.java

@ -15,7 +15,6 @@
*/ */
package org.thingsboard.server.queue.provider; 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.common.data.queue.Queue;
import org.thingsboard.server.gen.js.JsInvokeProtos; import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; 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 * Used to consume messages by TB Rule Engine Service
* *
* @return
* @param configuration * @param configuration
* @return
*/ */
TbQueueConsumer<TbProtoQueueMsg<ToRuleEngineMsg>> createToRuleEngineMsgConsumer(Queue configuration); TbQueueConsumer<TbProtoQueueMsg<ToRuleEngineMsg>> createToRuleEngineMsgConsumer(Queue configuration);
@ -105,9 +104,9 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory
* Used to consume messages by TB Rule Engine Service * Used to consume messages by TB Rule Engine Service
* Intended usage for consumer per partition strategy * Intended usage for consumer per partition strategy
* *
* @return TbQueueConsumer
* @param configuration * @param configuration
* @param partitionId as a suffix for consumer name * @param partitionId as a suffix for consumer name
* @return TbQueueConsumer
*/ */
default TbQueueConsumer<TbProtoQueueMsg<ToRuleEngineMsg>> createToRuleEngineMsgConsumer(Queue configuration, Integer partitionId) { default TbQueueConsumer<TbProtoQueueMsg<ToRuleEngineMsg>> createToRuleEngineMsgConsumer(Queue configuration, Integer partitionId) {
return createToRuleEngineMsgConsumer(configuration); return createToRuleEngineMsgConsumer(configuration);
@ -122,7 +121,7 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory
TbQueueRequestTemplate<TbProtoJsQueueMsg<JsInvokeProtos.RemoteJsRequest>, TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> createRemoteJsRequestTemplate(); TbQueueRequestTemplate<TbProtoJsQueueMsg<JsInvokeProtos.RemoteJsRequest>, TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> createRemoteJsRequestTemplate();
TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> createToCalculatedFieldMsgConsumer(TenantId tenantId); TbQueueConsumer<TbProtoQueueMsg<ToCalculatedFieldMsg>> createToCalculatedFieldMsgConsumer(Queue queue, Integer partitionId);
TbQueueAdmin getCalculatedFieldQueueAdmin(); TbQueueAdmin getCalculatedFieldQueueAdmin();
@ -132,7 +131,7 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory
TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> createToCalculatedFieldNotificationMsgProducer(); TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> createToCalculatedFieldNotificationMsgProducer();
TbQueueConsumer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateConsumer(TenantId tenantId); TbQueueConsumer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateConsumer();
TbQueueProducer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateProducer(); TbQueueProducer<TbProtoQueueMsg<CalculatedFieldStateProto>> createCalculatedFieldStateProducer();

Loading…
Cancel
Save