From 8519387ca7234e29175f7eb8750a7804adb22ff1 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 7 May 2024 00:14:15 +0200 Subject: [PATCH] Rule engine: Kafka consumer group per partition --- .../TbRuleEngineQueueConsumerManager.java | 7 ++++--- .../server/queue/discovery/TopicService.java | 7 ++++++- .../queue/provider/KafkaMonolithQueueFactory.java | 11 ++++++++++- .../provider/KafkaTbRuleEngineQueueFactory.java | 10 +++++++++- .../queue/provider/TbRuleEngineQueueFactory.java | 14 +++++++++++++- 5 files changed, 42 insertions(+), 7 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java index 2da3dbc6dc..38e54e8873 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java @@ -438,8 +438,9 @@ public class TbRuleEngineQueueConsumerManager { }); addedPartitions.forEach((tpi) -> { - String key = queueKey + "-" + tpi.getPartition().orElse(-999999); - TbQueueConsumerTask consumer = new TbQueueConsumerTask(key, ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue)); + int partitionId = tpi.getPartition().orElse(-999999); + String key = queueKey + "-" + partitionId; + TbQueueConsumerTask consumer = new TbQueueConsumerTask(key, ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, partitionId)); consumers.put(tpi, consumer); consumer.subscribe(Set.of(tpi)); launchConsumer(consumer); @@ -468,7 +469,7 @@ public class TbRuleEngineQueueConsumerManager { } if (consumer == null) { - consumer = new TbQueueConsumerTask(queueKey, ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue)); + consumer = new TbQueueConsumerTask(queueKey, ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, null)); } consumer.subscribe(partitions); if (!consumer.isRunning()) { 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 46084f8201..b7e000197e 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 @@ -64,4 +64,9 @@ public class TopicService { public String buildTopicName(String topic) { return prefix.isBlank() ? topic : prefix + "." + topic; } -} \ No newline at end of file + + public String suffix(Integer partitionId) { + return partitionId == null ? "" : "-" + partitionId; + } + +} 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 cd09f682c6..73c62e33dc 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 @@ -187,18 +187,27 @@ public class KafkaMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngi @Override public TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration) { + throw new UnsupportedOperationException("Rule engine consumer should use a partitionId"); + } + + @Override + public TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration, Integer partitionId) { String queueName = configuration.getName(); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(topicService.buildTopicName(configuration.getTopic())); consumerBuilder.clientId("re-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); - consumerBuilder.groupId(topicService.buildTopicName("re-" + queueName + (configuration.getTenantId().isSysTenantId() ? "" : ("-isolated-" + configuration.getTenantId())) + "-consumer")); + consumerBuilder.groupId(topicService.buildTopicName("re-" + queueName + + (configuration.getTenantId().isSysTenantId() ? "" : ("-isolated-" + configuration.getTenantId())) + + "-consumer" + + topicService.suffix(partitionId))); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(ruleEngineAdmin); consumerBuilder.statsService(consumerStatsService); return consumerBuilder.build(); } + @Override public TbQueueConsumer> createToRuleEngineNotificationsMsgConsumer() { TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueFactory.java index 568380d7b5..36238c206c 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 @@ -164,12 +164,20 @@ public class KafkaTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory { @Override public TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration) { + throw new UnsupportedOperationException("Rule engine consumer should use a partitionId"); + } + + @Override + public TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration, Integer partitionId) { String queueName = configuration.getName(); TbKafkaConsumerTemplate.TbKafkaConsumerTemplateBuilder> consumerBuilder = TbKafkaConsumerTemplate.builder(); consumerBuilder.settings(kafkaSettings); consumerBuilder.topic(topicService.buildTopicName(configuration.getTopic())); consumerBuilder.clientId("re-" + queueName + "-consumer-" + serviceInfoProvider.getServiceId() + "-" + consumerCount.incrementAndGet()); - consumerBuilder.groupId(topicService.buildTopicName("re-" + queueName + (configuration.getTenantId().isSysTenantId() ? "" : ("-isolated-" + configuration.getTenantId())) + "-consumer")); + consumerBuilder.groupId(topicService.buildTopicName("re-" + queueName + + (configuration.getTenantId().isSysTenantId() ? "" : ("-isolated-" + configuration.getTenantId())) + + "-consumer" + + topicService.suffix(partitionId))); consumerBuilder.decoder(msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders())); consumerBuilder.admin(ruleEngineAdmin); consumerBuilder.statsService(consumerStatsService); 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 f23c7e47f3..e49074f92d 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 @@ -77,13 +77,25 @@ public interface TbRuleEngineQueueFactory extends TbUsageStatsClientQueueFactory TbQueueProducer> createToOtaPackageStateServiceMsgProducer(); /** - * Used to consume messages by TB Core Service + * Used to consume messages by TB Rule Engine Service * * @return * @param configuration */ TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration); + /** + * 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 + */ + default TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration, Integer partitionId) { + return createToRuleEngineMsgConsumer(configuration); + } + /** * Used to consume high priority messages by TB Core Service *