Browse Source

Fix cluster service send to re

pull/14074/head
Andrii Landiak 11 months ago
parent
commit
a39e8f79b2
  1. 17
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

17
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -627,15 +627,14 @@ public class DefaultTbClusterService implements TbClusterService {
// No need to push notifications twice // No need to push notifications twice
tbRuleEngineServices.removeAll(tbCoreServices); tbRuleEngineServices.removeAll(tbCoreServices);
} }
if (entityType == EntityType.USER) { boolean toRuleEngine = entityType != EntityType.USER;
// No need to push user update notification to the rule engine if (toRuleEngine) {
return; for (String serviceId : tbRuleEngineServices) {
} TopicPartitionInfo tpi = topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceId);
for (String serviceId : tbRuleEngineServices) { ToRuleEngineNotificationMsg toRuleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setComponentLifecycle(componentLifecycleMsgProto).build();
TopicPartitionInfo tpi = topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceId); toRuleEngineProducer.send(tpi, new TbProtoQueueMsg<>(msg.getEntityId().getId(), toRuleEngineMsg), null);
ToRuleEngineNotificationMsg toRuleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setComponentLifecycle(componentLifecycleMsgProto).build(); toRuleEngineNfs.incrementAndGet();
toRuleEngineProducer.send(tpi, new TbProtoQueueMsg<>(msg.getEntityId().getId(), toRuleEngineMsg), null); }
toRuleEngineNfs.incrementAndGet();
} }
} }

Loading…
Cancel
Save