From a5404c2b4576021a2b5b061347216256a2b23188 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Mon, 24 Mar 2025 17:26:01 +0200 Subject: [PATCH] EDQS: don't use consumer group management for state topics; single consumer group for requests topic; reduce events topic retention to 24 hours --- .../src/main/resources/thingsboard.yml | 2 +- .../edqs/state/KafkaEdqsStateService.java | 4 +-- .../server/queue/discovery/TopicService.java | 9 ++++-- .../queue/discovery/ZkDiscoveryService.java | 2 +- .../queue/edqs/KafkaEdqsQueueFactory.java | 3 +- .../server/queue/kafka/TbKafkaAdmin.java | 2 ++ .../queue/kafka/TbKafkaConsumerTemplate.java | 28 ++++++++++++------- edqs/src/main/resources/edqs.yml | 2 +- 8 files changed, 31 insertions(+), 21 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index eac787d52b..7737f9d027 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1646,7 +1646,7 @@ queue: # Kafka properties for Calculated Field State topics calculated-field-state: "${TB_QUEUE_KAFKA_CF_STATE_TOPIC_PROPERTIES:retention.ms:-1;segment.bytes:52428800;retention.bytes:104857600000;partitions:1;min.insync.replicas:1;cleanup.policy:compact}" # Kafka properties for EDQS events topics - edqs-events: "${TB_QUEUE_KAFKA_EDQS_EVENTS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:52428800;retention.bytes:-1;partitions:1;min.insync.replicas:1}" + edqs-events: "${TB_QUEUE_KAFKA_EDQS_EVENTS_TOPIC_PROPERTIES:retention.ms:86400000;segment.bytes:52428800;retention.bytes:-1;partitions:1;min.insync.replicas:1}" # Kafka properties for EDQS requests topic (default: 3 minutes retention) edqs-requests: "${TB_QUEUE_KAFKA_EDQS_REQUESTS_TOPIC_PROPERTIES:retention.ms:180000;segment.bytes:52428800;retention.bytes:1048576000;partitions:1;min.insync.replicas:1}" # Kafka properties for EDQS state topic (infinite retention, compaction) 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 c59707c9c3..21515f1ad8 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 @@ -36,7 +36,6 @@ import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; import org.thingsboard.server.queue.common.state.KafkaQueueStateService; import org.thingsboard.server.queue.common.state.QueueStateService; import org.thingsboard.server.queue.discovery.QueueKey; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.edqs.EdqsConfig; import org.thingsboard.server.queue.edqs.EdqsQueue; import org.thingsboard.server.queue.edqs.EdqsQueueFactory; @@ -57,7 +56,6 @@ public class KafkaEdqsStateService implements EdqsStateService { private final EdqsConfig config; private final EdqsPartitionService partitionService; private final EdqsQueueFactory queueFactory; - private final TopicService topicService; @Autowired @Lazy private EdqsProcessor edqsProcessor; @@ -91,7 +89,7 @@ public class KafkaEdqsStateService implements EdqsStateService { } consumer.commit(); }) - .consumerCreator((config, partitionId) -> queueFactory.createEdqsMsgConsumer(EdqsQueue.STATE)) + .consumerCreator((config, partitionId) -> queueFactory.createEdqsMsgConsumer(EdqsQueue.STATE, null)) // not using consumer group management .queueAdmin(queueFactory.getEdqsQueueAdmin()) .consumerExecutor(eventConsumer.getConsumerExecutor()) .taskExecutor(eventConsumer.getTaskExecutor()) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java index 5992083d85..2a0fad0643 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java @@ -103,6 +103,9 @@ public class TopicService { } public String buildTopicName(String topic) { + if (topic == null) { + return null; + } return prefix.isBlank() ? topic : prefix + "." + topic; } @@ -113,9 +116,9 @@ public class TopicService { public String buildConsumerGroupId(String servicePrefix, TenantId tenantId, String queueName, Integer partitionId) { return this.buildTopicName( servicePrefix + queueName - + (tenantId.isSysTenantId() ? "" : ("-isolated-" + tenantId)) - + "-consumer" - + suffix(partitionId)); + + (tenantId.isSysTenantId() ? "" : ("-isolated-" + tenantId)) + + "-consumer" + + suffix(partitionId)); } String suffix(Integer partitionId) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java index f7a4d2abf6..cf9f27ee39 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java @@ -315,7 +315,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi ScheduledFuture task = delayedTasks.remove(serviceId); if (task != null) { if (task.cancel(false)) { - log.debug("[{}] Recalculate partitions ignored. Service was restarted in time [{}].", + log.info("[{}] Recalculate partitions ignored. Service was restarted in time [{}].", serviceId, serviceTypesList); } else { log.debug("[{}] Going to recalculate partitions. Service was not restarted in time [{}]!", diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java index a322cc5434..6fdab28133 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java @@ -103,12 +103,11 @@ public class KafkaEdqsQueueFactory implements EdqsQueueFactory { @Override public TbQueueResponseTemplate, TbProtoQueueMsg> createEdqsResponseTemplate() { - String requestsConsumerGroup = "edqs-requests-consumer-group-" + edqsConfig.getLabel(); var requestConsumer = TbKafkaConsumerTemplate.>builder() .settings(kafkaSettings) .topic(topicService.buildTopicName(edqsConfig.getRequestsTopic())) .clientId("edqs-requests-consumer-" + serviceInfoProvider.getServiceId()) - .groupId(topicService.buildTopicName(requestsConsumerGroup)) + .groupId(topicService.buildTopicName("edqs-requests-consumer-group")) .decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportProtos.ToEdqsMsg.parseFrom(msg.getData()), msg.getHeaders())) .admin(edqsRequestsAdmin) .statsService(consumerStatsService); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java index 3496aac76a..6835b44da7 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.queue.kafka; +import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.admin.CreateTopicsResult; import org.apache.kafka.clients.admin.ListOffsetsResult; @@ -45,6 +46,7 @@ public class TbKafkaAdmin implements TbQueueAdmin { private final TbKafkaSettings settings; private final Map topicConfigs; + @Getter private final int numPartitions; private volatile Set topics; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java index 4bd3bf0fe6..8c4ea788c6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java @@ -40,6 +40,7 @@ import java.util.Map; import java.util.Properties; import java.util.Set; import java.util.stream.Collectors; +import java.util.stream.IntStream; /** * Created by ashvayka on 24.09.18. @@ -47,7 +48,7 @@ import java.util.stream.Collectors; @Slf4j public class TbKafkaConsumerTemplate extends AbstractTbQueueConsumerTemplate, T> { - private final TbQueueAdmin admin; + private final TbKafkaAdmin admin; private final KafkaConsumer consumer; private final TbKafkaDecoder decoder; @@ -78,7 +79,7 @@ public class TbKafkaConsumerTemplate extends AbstractTbQue statsService.registerClientGroup(groupId); } - this.admin = admin; + this.admin = (TbKafkaAdmin) admin; this.consumer = new KafkaConsumer<>(props); this.decoder = decoder; this.readFromBeginning = readFromBeginning; @@ -105,14 +106,19 @@ public class TbKafkaConsumerTemplate extends AbstractTbQue List toSubscribe = new ArrayList<>(); topics.forEach((topic, kafkaPartitions) -> { if (kafkaPartitions == null) { - toSubscribe.add(topic); - } else { - List topicPartitions = kafkaPartitions.stream() - .map(partition -> new TopicPartition(topic, partition)) - .toList(); - consumer.assign(topicPartitions); - onPartitionsAssigned(topicPartitions); + if (groupId != null) { + toSubscribe.add(topic); + return; + } else { // if no consumer group management - manually assigning all topic partitions + kafkaPartitions = IntStream.range(0, admin.getNumPartitions()).boxed().toList(); + } } + + List topicPartitions = kafkaPartitions.stream() + .map(partition -> new TopicPartition(topic, partition)) + .toList(); + consumer.assign(topicPartitions); + onPartitionsAssigned(topicPartitions); }); if (!toSubscribe.isEmpty()) { if (readFromBeginning || stopWhenRead) { @@ -195,7 +201,9 @@ public class TbKafkaConsumerTemplate extends AbstractTbQue @Override protected void doCommit() { - consumer.commitSync(); + if (groupId != null) { + consumer.commitSync(); + } } @Override diff --git a/edqs/src/main/resources/edqs.yml b/edqs/src/main/resources/edqs.yml index c101eff68e..05d942ff23 100644 --- a/edqs/src/main/resources/edqs.yml +++ b/edqs/src/main/resources/edqs.yml @@ -149,7 +149,7 @@ queue: # value: "${TB_QUEUE_KAFKA_SESSION_TIMEOUT_MS:10000}" # (10 seconds) topic-properties: # Kafka properties for EDQS events topics - edqs-events: "${TB_QUEUE_KAFKA_EDQS_EVENTS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:52428800;retention.bytes:-1;partitions:1;min.insync.replicas:1}" + edqs-events: "${TB_QUEUE_KAFKA_EDQS_EVENTS_TOPIC_PROPERTIES:retention.ms:86400000;segment.bytes:52428800;retention.bytes:-1;partitions:1;min.insync.replicas:1}" # Kafka properties for EDQS requests topic (default: 3 minutes retention) edqs-requests: "${TB_QUEUE_KAFKA_EDQS_REQUESTS_TOPIC_PROPERTIES:retention.ms:180000;segment.bytes:52428800;retention.bytes:1048576000;partitions:1;min.insync.replicas:1}" # Kafka properties for EDQS state topic (infinite retention, compaction)