From f32e2f6fdefb96cef3b8dc13008e76326b4faf00 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 10 Aug 2023 12:03:15 +0300 Subject: [PATCH] Refactoring after review --- .../service/entitiy/queue/DefaultTbQueueService.java | 4 ++++ .../queue/DefaultTbRuleEngineConsumerService.java | 6 +++--- .../org/thingsboard/server/queue/TbQueueConsumer.java | 2 +- .../queue/common/AbstractTbQueueConsumerTemplate.java | 9 ++++----- .../server/queue/memory/InMemoryTbQueueConsumer.java | 8 ++++---- 5 files changed, 16 insertions(+), 13 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java index 72b8c2252a..9aa9a421c8 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java @@ -176,6 +176,10 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb } } + if (log.isDebugEnabled()) { + log.debug("[{}] Handling profile queue config update: creating queues {}, updating {}, deleting {}. Affected tenants: {}", + newTenantProfile.getUuidId(), toCreate, toUpdate, toRemove, tenantIds); + } tenantIds.forEach(tenantId -> { toCreate.forEach(key -> saveQueue(new Queue(tenantId, newQueues.get(key)))); 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 863f0354c2..c8bf469a61 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 @@ -121,7 +121,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< private final ConcurrentMap consumerStats = new ConcurrentHashMap<>(); private final ConcurrentMap topicsConsumerPerPartition = new ConcurrentHashMap<>(); final ExecutorService submitExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("tb-rule-engine-consumer-submit")); - final ScheduledExecutorService repartitionExecutor = Executors.newScheduledThreadPool(2, ThingsBoardThreadFactory.forName("tb-rule-engine-consumer-repartition")); + final ScheduledExecutorService repartitionExecutor = Executors.newScheduledThreadPool(1, ThingsBoardThreadFactory.forName("tb-rule-engine-consumer-repartition")); public DefaultTbRuleEngineConsumerService(TbRuleEngineProcessingStrategyFactory processingStrategyFactory, TbRuleEngineSubmitStrategyFactory submitStrategyFactory, @@ -274,7 +274,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< void consumerLoop(TbQueueConsumer> consumer, org.thingsboard.server.common.data.queue.Queue configuration, TbRuleEngineConsumerStats stats, String threadSuffix) { updateCurrentThreadName(threadSuffix); - while (!stopped && !consumer.isStopped() && !consumer.isDeleted()) { + while (!stopped && !consumer.isStopped() && !consumer.isQueueDeleted()) { try { List> msgs = consumer.poll(configuration.getPollInterval()); if (msgs.isEmpty()) { @@ -325,7 +325,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< } } - if (consumer.isDeleted()) { + if (consumer.isQueueDeleted()) { processQueueDeletion(configuration, consumer); } log.info("TB Rule Engine Consumer stopped."); diff --git a/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java index 73bea7642f..9c41f9d342 100644 --- a/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java +++ b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java @@ -38,7 +38,7 @@ public interface TbQueueConsumer { void onQueueDelete(); - boolean isDeleted(); + boolean isQueueDeleted(); List getFullTopicNames(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueConsumerTemplate.java index 8e91c1d134..2ebe41850d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueConsumerTemplate.java @@ -44,7 +44,7 @@ public abstract class AbstractTbQueueConsumerTemplate i protected volatile Set partitions; protected final ReentrantLock consumerLock = new ReentrantLock(); //NonfairSync final Queue> subscribeQueue = new ConcurrentLinkedQueue<>(); - protected volatile boolean deleted = false; + protected volatile boolean queueDeleted = false; @Getter private final String topic; @@ -194,12 +194,11 @@ public abstract class AbstractTbQueueConsumerTemplate i @Override public void onQueueDelete() { - deleted = true; + queueDeleted = true; } - @Override - public boolean isDeleted() { - return deleted; + public boolean isQueueDeleted() { + return queueDeleted; } @Override diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java index dba6f6d588..8711cbbcf1 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java @@ -31,7 +31,7 @@ public class InMemoryTbQueueConsumer implements TbQueueCon private volatile Set partitions; private volatile boolean stopped; private volatile boolean subscribed; - private volatile boolean deleted; + private volatile boolean queueDeleted; public InMemoryTbQueueConsumer(InMemoryStorage storage, String topic) { this.storage = storage; @@ -106,12 +106,12 @@ public class InMemoryTbQueueConsumer implements TbQueueCon @Override public void onQueueDelete() { - deleted = true; + queueDeleted = true; } @Override - public boolean isDeleted() { - return deleted; + public boolean isQueueDeleted() { + return queueDeleted; } @Override