Browse Source

Refactoring after review

pull/8988/head
ViacheslavKlimov 3 years ago
parent
commit
f32e2f6fde
  1. 4
      application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java
  2. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  3. 2
      common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java
  4. 9
      common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueConsumerTemplate.java
  5. 8
      common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java

4
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 -> { tenantIds.forEach(tenantId -> {
toCreate.forEach(key -> saveQueue(new Queue(tenantId, newQueues.get(key)))); toCreate.forEach(key -> saveQueue(new Queue(tenantId, newQueues.get(key))));

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -121,7 +121,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
private final ConcurrentMap<QueueKey, TbRuleEngineConsumerStats> consumerStats = new ConcurrentHashMap<>(); private final ConcurrentMap<QueueKey, TbRuleEngineConsumerStats> consumerStats = new ConcurrentHashMap<>();
private final ConcurrentMap<QueueKey, TbTopicWithConsumerPerPartition> topicsConsumerPerPartition = new ConcurrentHashMap<>(); private final ConcurrentMap<QueueKey, TbTopicWithConsumerPerPartition> topicsConsumerPerPartition = new ConcurrentHashMap<>();
final ExecutorService submitExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("tb-rule-engine-consumer-submit")); 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, public DefaultTbRuleEngineConsumerService(TbRuleEngineProcessingStrategyFactory processingStrategyFactory,
TbRuleEngineSubmitStrategyFactory submitStrategyFactory, TbRuleEngineSubmitStrategyFactory submitStrategyFactory,
@ -274,7 +274,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
void consumerLoop(TbQueueConsumer<TbProtoQueueMsg<ToRuleEngineMsg>> consumer, org.thingsboard.server.common.data.queue.Queue configuration, TbRuleEngineConsumerStats stats, String threadSuffix) { void consumerLoop(TbQueueConsumer<TbProtoQueueMsg<ToRuleEngineMsg>> consumer, org.thingsboard.server.common.data.queue.Queue configuration, TbRuleEngineConsumerStats stats, String threadSuffix) {
updateCurrentThreadName(threadSuffix); updateCurrentThreadName(threadSuffix);
while (!stopped && !consumer.isStopped() && !consumer.isDeleted()) { while (!stopped && !consumer.isStopped() && !consumer.isQueueDeleted()) {
try { try {
List<TbProtoQueueMsg<ToRuleEngineMsg>> msgs = consumer.poll(configuration.getPollInterval()); List<TbProtoQueueMsg<ToRuleEngineMsg>> msgs = consumer.poll(configuration.getPollInterval());
if (msgs.isEmpty()) { if (msgs.isEmpty()) {
@ -325,7 +325,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
} }
} }
if (consumer.isDeleted()) { if (consumer.isQueueDeleted()) {
processQueueDeletion(configuration, consumer); processQueueDeletion(configuration, consumer);
} }
log.info("TB Rule Engine Consumer stopped."); log.info("TB Rule Engine Consumer stopped.");

2
common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java

@ -38,7 +38,7 @@ public interface TbQueueConsumer<T extends TbQueueMsg> {
void onQueueDelete(); void onQueueDelete();
boolean isDeleted(); boolean isQueueDeleted();
List<String> getFullTopicNames(); List<String> getFullTopicNames();

9
common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueConsumerTemplate.java

@ -44,7 +44,7 @@ public abstract class AbstractTbQueueConsumerTemplate<R, T extends TbQueueMsg> i
protected volatile Set<TopicPartitionInfo> partitions; protected volatile Set<TopicPartitionInfo> partitions;
protected final ReentrantLock consumerLock = new ReentrantLock(); //NonfairSync protected final ReentrantLock consumerLock = new ReentrantLock(); //NonfairSync
final Queue<Set<TopicPartitionInfo>> subscribeQueue = new ConcurrentLinkedQueue<>(); final Queue<Set<TopicPartitionInfo>> subscribeQueue = new ConcurrentLinkedQueue<>();
protected volatile boolean deleted = false; protected volatile boolean queueDeleted = false;
@Getter @Getter
private final String topic; private final String topic;
@ -194,12 +194,11 @@ public abstract class AbstractTbQueueConsumerTemplate<R, T extends TbQueueMsg> i
@Override @Override
public void onQueueDelete() { public void onQueueDelete() {
deleted = true; queueDeleted = true;
} }
@Override public boolean isQueueDeleted() {
public boolean isDeleted() { return queueDeleted;
return deleted;
} }
@Override @Override

8
common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java

@ -31,7 +31,7 @@ public class InMemoryTbQueueConsumer<T extends TbQueueMsg> implements TbQueueCon
private volatile Set<TopicPartitionInfo> partitions; private volatile Set<TopicPartitionInfo> partitions;
private volatile boolean stopped; private volatile boolean stopped;
private volatile boolean subscribed; private volatile boolean subscribed;
private volatile boolean deleted; private volatile boolean queueDeleted;
public InMemoryTbQueueConsumer(InMemoryStorage storage, String topic) { public InMemoryTbQueueConsumer(InMemoryStorage storage, String topic) {
this.storage = storage; this.storage = storage;
@ -106,12 +106,12 @@ public class InMemoryTbQueueConsumer<T extends TbQueueMsg> implements TbQueueCon
@Override @Override
public void onQueueDelete() { public void onQueueDelete() {
deleted = true; queueDeleted = true;
} }
@Override @Override
public boolean isDeleted() { public boolean isQueueDeleted() {
return deleted; return queueDeleted;
} }
@Override @Override

Loading…
Cancel
Save