From 12af0cd757821915de4a2969698bd354976f40a5 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 13 Feb 2024 15:01:15 +0200 Subject: [PATCH] Always recalculate partitions after queue updates --- .../DefaultTbRuleEngineConsumerService.java | 20 +++++-------------- 1 file changed, 5 insertions(+), 15 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java index b92ed34639..afdd150a59 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java @@ -184,7 +184,6 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< } private void updateQueues(List queueUpdateMsgs) { - boolean partitionsChanged = false; for (QueueUpdateMsg queueUpdateMsg : queueUpdateMsgs) { log.info("Received queue update msg: [{}]", queueUpdateMsg); TenantId tenantId = new TenantId(new UUID(queueUpdateMsg.getTenantIdMSB(), queueUpdateMsg.getTenantIdLSB())); @@ -194,23 +193,14 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueName, tenantId); Queue queue = queueService.findQueueById(tenantId, queueId); - var consumer = getConsumer(queueKey).orElseGet(() -> createConsumer(queueKey, queue)); - Queue oldQueue = consumer.getQueue(); - consumer.update(queue); - - if (oldQueue == null || queue.getPartitions() != oldQueue.getPartitions()) { - partitionsChanged = true; - } - } else { - partitionsChanged = true; + getConsumer(queueKey).ifPresentOrElse(consumer -> consumer.update(queue), + () -> createConsumer(queueKey, queue)); } } - if (partitionsChanged) { - partitionService.updateQueues(queueUpdateMsgs); - partitionService.recalculatePartitions(ctx.getServiceInfoProvider().getServiceInfo(), - new ArrayList<>(partitionService.getOtherServices(ServiceType.TB_RULE_ENGINE))); - } + partitionService.updateQueues(queueUpdateMsgs); + partitionService.recalculatePartitions(ctx.getServiceInfoProvider().getServiceInfo(), + new ArrayList<>(partitionService.getOtherServices(ServiceType.TB_RULE_ENGINE))); } private void deleteQueues(List queueDeleteMsgs) {