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 bc28042fb2..9ac8f86590 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 @@ -22,8 +22,10 @@ import akka.actor.Terminated; import com.google.protobuf.ByteString; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.context.event.EventListener; +import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.thingsboard.rule.engine.api.msg.DeviceCredentialsUpdateNotificationMsg; import org.thingsboard.rule.engine.api.msg.DeviceNameOrTypeUpdateMsg; @@ -50,6 +52,7 @@ import org.thingsboard.server.service.cluster.discovery.DiscoveryService; import org.thingsboard.server.service.cluster.discovery.ServerInstance; import org.thingsboard.server.service.cluster.rpc.ClusterRpcService; import org.thingsboard.server.service.state.DeviceStateService; +import org.thingsboard.server.service.transport.RuleEngineStats; import scala.concurrent.Await; import scala.concurrent.Future; import scala.concurrent.duration.Duration; @@ -58,6 +61,7 @@ import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.atomic.AtomicInteger; import static org.thingsboard.server.gen.cluster.ClusterAPIProtos.MessageType.CLUSTER_ACTOR_MESSAGE; @@ -187,8 +191,25 @@ public class DefaultActorService implements ActorService { this.rpcManagerActor.tell(msg, ActorRef.noSender()); } + @Value("${cluster.stats.enabled:false}") + private boolean statsEnabled; + + private final AtomicInteger sentClusterMsgs = new AtomicInteger(0); + private final AtomicInteger receivedClusterMsgs = new AtomicInteger(0); + + + @Scheduled(fixedDelayString = "${cluster.stats.print_interval_ms}") + public void printStats() { + if (statsEnabled) { + log.info("Cluster msgs sent [{}] received [{}]", sentClusterMsgs.getAndSet(0), receivedClusterMsgs.getAndSet(0)); + } + } + @Override public void onReceivedMsg(ServerAddress source, ClusterAPIProtos.ClusterMessage msg) { + if (statsEnabled) { + receivedClusterMsgs.incrementAndGet(); + } ServerAddress serverAddress = new ServerAddress(source.getHost(), source.getPort(), source.getServerType()); if (log.isDebugEnabled()) { log.info("Received msg [{}] from [{}]", msg.getMessageType().name(), serverAddress); @@ -239,11 +260,17 @@ public class DefaultActorService implements ActorService { @Override public void onSendMsg(ClusterAPIProtos.ClusterMessage msg) { + if (statsEnabled) { + sentClusterMsgs.incrementAndGet(); + } rpcManagerActor.tell(msg, ActorRef.noSender()); } @Override public void onRpcSessionCreateRequestMsg(RpcSessionCreateRequestMsg msg) { + if (statsEnabled) { + sentClusterMsgs.incrementAndGet(); + } rpcManagerActor.tell(msg, ActorRef.noSender()); } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index b181572247..aa10b6c097 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -85,6 +85,9 @@ cluster: vitrual_nodes_size: "${CLUSTER_VIRTUAL_NODES_SIZE:16}" # Queue partition id for current node partition_id: "${QUEUE_PARTITION_ID:0}" + stats: + enabled: "${TB_RULE_ENGINE_STATS_ENABLED:false}" + print_interval_ms: "${TB_RULE_ENGINE_STATS_PRINT_INTERVAL_MS:10000}" # Plugins configuration parameters plugins: