From 416439fab08bd59997bfa93a6457fe7862575cee Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Mon, 13 May 2024 17:19:58 +0200 Subject: [PATCH] TbRuleEngineQueueConsumerManager - createToRuleEngineMsgConsumer lazy --- .../queue/ruleengine/TbQueueConsumerTask.java | 27 ++++++++++++++++--- .../TbRuleEngineQueueConsumerManager.java | 7 ++--- 2 files changed, 26 insertions(+), 8 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbQueueConsumerTask.java b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbQueueConsumerTask.java index 0e25efd014..a25f74a5b1 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbQueueConsumerTask.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbQueueConsumerTask.java @@ -24,22 +24,43 @@ import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import java.util.Objects; import java.util.Set; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; -@RequiredArgsConstructor @Slf4j public class TbQueueConsumerTask { @Getter private final Object key; - @Getter - private final TbQueueConsumer> consumer; + private volatile TbQueueConsumer> consumer; + private volatile Supplier>> consumerSupplier; @Setter private Future task; + public TbQueueConsumerTask(Object key, Supplier>> consumerSupplier) { + this.key = key; + this.consumer = null; + this.consumerSupplier = consumerSupplier; + } + + public TbQueueConsumer> getConsumer() { + if (consumer == null) { + synchronized (this) { + if (consumer == null) { + Objects.requireNonNull(consumerSupplier, "consumerSupplier for key [" + key + "] is null"); + consumer = consumerSupplier.get(); + Objects.requireNonNull(consumer, "consumer for key [" + key + "] is null"); + consumerSupplier = null; + } + } + } + return consumer; + } + public void subscribe(Set partitions) { log.trace("[{}] Subscribing to partitions: {}", key, partitions); consumer.subscribe(partitions); 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 39d91341f0..4df0a900d2 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 @@ -433,16 +433,13 @@ public class TbRuleEngineQueueConsumerManager { removedPartitions.forEach((tpi) -> { consumers.get(tpi).initiateStop(); }); - removedPartitions.forEach((tpi) -> { - consumers.remove(tpi).awaitCompletion(); - }); // ctx.getQueueAdmin().createTopicIfNotExists(); addedPartitions.forEach((tpi) -> { int partitionId = tpi.getPartition().orElse(-999999); String key = queueKey + "-" + partitionId; - TbQueueConsumerTask consumer = new TbQueueConsumerTask(key, ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, null)); + TbQueueConsumerTask consumer = new TbQueueConsumerTask(key, () -> ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, partitionId)); consumers.put(tpi, consumer); consumer.subscribe(Set.of(tpi)); launchConsumer(consumer); @@ -471,7 +468,7 @@ public class TbRuleEngineQueueConsumerManager { } if (consumer == null) { - consumer = new TbQueueConsumerTask(queueKey, ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, null)); + consumer = new TbQueueConsumerTask(queueKey, () -> ctx.getQueueFactory().createToRuleEngineMsgConsumer(queue, null)); } consumer.subscribe(partitions); if (!consumer.isRunning()) {