From c1f16769f47a1d606a16ba61769ae30993660084 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 9 Feb 2024 14:15:24 +0200 Subject: [PATCH] Fix partition change event ignored --- .../actors/service/DefaultActorService.java | 6 ++++++ .../AbstractPartitionBasedService.java | 13 ++++++++----- .../queue/DefaultTbCoreConsumerService.java | 18 ++++++++---------- .../DefaultTbRuleEngineConsumerService.java | 18 ++++++++---------- .../processing/AbstractConsumerService.java | 5 +++++ .../discovery/TbApplicationEventListener.java | 12 +++++++++--- 6 files changed, 44 insertions(+), 28 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java index 87831cb8f6..0ecdb183b5 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java @@ -31,6 +31,7 @@ import org.thingsboard.server.actors.app.AppActor; import org.thingsboard.server.actors.app.AppInitMsg; import org.thingsboard.server.actors.stats.StatsActor; import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; +import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.util.AfterStartUp; @@ -124,6 +125,11 @@ public class DefaultActorService extends TbApplicationEventListener extends */ @Override protected void onTbApplicationEvent(PartitionChangeEvent partitionChangeEvent) { - if (getServiceType().equals(partitionChangeEvent.getServiceType())) { - log.debug("onTbApplicationEvent, processing event: {}", partitionChangeEvent); - subscribeQueue.add(partitionChangeEvent.getPartitions()); - scheduledExecutor.submit(this::pollInitStateFromDB); - } + log.debug("onTbApplicationEvent, processing event: {}", partitionChangeEvent); + subscribeQueue.add(partitionChangeEvent.getPartitions()); + scheduledExecutor.submit(this::pollInitStateFromDB); + } + + @Override + protected boolean filterTbApplicationEvent(PartitionChangeEvent event) { + return getServiceType().equals(event.getServiceType()); } protected void pollInitStateFromDB() { diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index f544c4fe25..f846d8ce28 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -220,16 +220,14 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService tpi.newByTopic(usageStatsConsumer.getTopic())) - .collect(Collectors.toSet())); - } + log.info("Subscribing to partitions: {}", event.getPartitions()); + this.mainConsumer.subscribe(event.getPartitions()); + this.usageStatsConsumer.subscribe( + event + .getPartitions() + .stream() + .map(tpi -> tpi.newByTopic(usageStatsConsumer.getTopic())) + .collect(Collectors.toSet())); this.firmwareStatesConsumer.subscribe(); } 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 6a205ed5a2..236b85e205 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 @@ -110,16 +110,14 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< @Override protected void onTbApplicationEvent(PartitionChangeEvent event) { - if (event.getServiceType().equals(getServiceType())) { - event.getPartitionsMap().forEach((queueKey, partitions) -> { - var consumer = consumers.get(queueKey); - if (consumer != null) { - consumer.update(partitions); - } else { - log.warn("Received invalid partition change event for {} that is not managed by this service", queueKey); - } - }); - } + event.getPartitionsMap().forEach((queueKey, partitions) -> { + var consumer = consumers.get(queueKey); + if (consumer != null) { + consumer.update(partitions); + } else { + log.warn("Received invalid partition change event for {} that is not managed by this service", queueKey); + } + }); } @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java index a8252e7fa0..8a728d6579 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java @@ -108,6 +108,11 @@ public abstract class AbstractConsumerService implements ApplicationListener { private int lastProcessedSequenceNumber = Integer.MIN_VALUE; private final Lock seqNumberLock = new ReentrantLock(); + private final Logger log = LoggerFactory.getLogger(getClass()); + @Override public void onApplicationEvent(T event) { + if (!filterTbApplicationEvent(event)) { + log.trace("Skipping event due to filter: {}", event); + return; + } boolean validUpdate = false; seqNumberLock.lock(); try { @@ -40,7 +46,7 @@ public abstract class TbApplicationEventListener i } finally { seqNumberLock.unlock(); } - if (validUpdate && filterTbApplicationEvent(event)) { + if (validUpdate) { try { onTbApplicationEvent(event); } catch (Exception e) {