From 8ede0b28226c97ad205886490138e0d966959386 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 15 Mar 2023 16:54:23 +0100 Subject: [PATCH] added ability to send system info throw ZK --- .../system/DefaultSystemInfoService.java | 116 ++++++++++++------ common/cluster-api/src/main/proto/queue.proto | 7 ++ .../server/common/data/SystemInfo.java | 12 +- .../server/common/data/SystemInfoData.java | 35 ++++++ .../DefaultTbServiceInfoProvider.java | 39 +++++- .../queue/discovery/DiscoveryService.java | 6 + .../discovery/DummyDiscoveryService.java | 6 + .../queue/discovery/HashPartitionService.java | 5 +- .../queue/discovery/PartitionService.java | 2 - .../discovery/TbServiceInfoProvider.java | 3 + .../queue/discovery/ZkDiscoveryService.java | 29 +++-- .../thingsboard/common/util/SystemUtil.java | 59 +++++++++ 12 files changed, 250 insertions(+), 69 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/SystemInfoData.java create mode 100644 common/util/src/main/java/org/thingsboard/common/util/SystemUtil.java diff --git a/application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java b/application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java index 3e9b778a5d..40ea65d9fd 100644 --- a/application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java +++ b/application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java @@ -15,77 +15,119 @@ */ package org.thingsboard.server.service.system; +import com.google.common.util.concurrent.FutureCallback; import com.google.protobuf.ProtocolStringList; import lombok.RequiredArgsConstructor; -import lombok.SneakyThrows; +import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; +import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.SystemInfo; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.common.data.SystemInfoData; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo; +import org.thingsboard.server.queue.discovery.DiscoveryService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; -import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; -import java.io.File; -import java.lang.management.ManagementFactory; -import java.lang.management.MemoryMXBean; -import java.lang.management.OperatingSystemMXBean; -import java.util.HashMap; +import javax.annotation.Nullable; +import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; +import java.util.ArrayList; +import java.util.Collections; import java.util.List; -import java.util.Map; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + +import static org.thingsboard.common.util.SystemUtil.getCpuUsage; +import static org.thingsboard.common.util.SystemUtil.getFreeDiscSpace; +import static org.thingsboard.common.util.SystemUtil.getMemoryUsage; -@TbCoreComponent @Service @RequiredArgsConstructor +@Slf4j public class DefaultSystemInfoService implements SystemInfoService { + public static final FutureCallback CALLBACK = new FutureCallback<>() { + @Override + public void onSuccess(@Nullable Integer result) { + } + + @Override + public void onFailure(Throwable t) { + log.warn("Failed to persist system info", t); + } + }; + private final TbServiceInfoProvider serviceInfoProvider; - private final PartitionService partitionService; + private final DiscoveryService discoveryService; + private final TelemetrySubscriptionService telemetryService; + private ScheduledExecutorService scheduler; @Value("${zk.enabled:false}") private boolean zkEnabled; + @PostConstruct + private void init() { + if (!zkEnabled) { + scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("tb-system-info-scheduler")); + scheduler.scheduleAtFixedRate(this::saveCurrentSystemInfo, 0, 1, TimeUnit.MINUTES); + } + } + + @PreDestroy + private void destroy() { + if (scheduler != null) { + scheduler.shutdownNow(); + } + } + @Override - @SneakyThrows public SystemInfo getSystemInfo() { SystemInfo systemInfo = new SystemInfo(); - TransportProtos.ServiceInfo serviceInfo = serviceInfoProvider.getServiceInfo(); - List currentOtherServices = partitionService.getCurrentOtherServices(); + ServiceInfo serviceInfo = serviceInfoProvider.getServiceInfoWithCurrentSystemInfo(); if (zkEnabled) { - Map serviceInfos = new HashMap<>(); - addServiceInfo(serviceInfos, serviceInfo); - currentOtherServices.forEach(otherInfo -> addServiceInfo(serviceInfos, otherInfo)); - systemInfo.setServiceInfos(serviceInfos); + List clusterSystemData = new ArrayList<>(); + clusterSystemData.add(createSystemInfoData(serviceInfo)); + this.discoveryService.getOtherServers() + .stream() + .map(this::createSystemInfoData) + .forEach(clusterSystemData::add); + systemInfo.setSystemData(clusterSystemData); } else { systemInfo.setMonolith(true); - systemInfo.setMemUsage(getMemoryUsage()); - - systemInfo.setCpuUsage((int) (getCpuUsage() * 100) / 100.0); - systemInfo.setFreeDiscSpace(getFreeDiscSpace()); + systemInfo.setSystemData(Collections.singletonList(createSystemInfoData(serviceInfo))); } return systemInfo; } - private void addServiceInfo(Map serviceInfos, TransportProtos.ServiceInfo serviceInfo) { - ProtocolStringList serviceTypes = serviceInfo.getServiceTypesList(); - serviceInfos.put(serviceInfo.getServiceId(), serviceTypes.size() > 1 ? "MONOLITH" : serviceTypes.get(0)); - } + private void saveCurrentSystemInfo() { + long ts = System.currentTimeMillis(); + List tsList = new ArrayList<>(); + tsList.add(new BasicTsKvEntry(ts, new LongDataEntry("memoryUsage", getMemoryUsage()))); + tsList.add(new BasicTsKvEntry(ts, new DoubleDataEntry("cpuUsage", getCpuUsage()))); + tsList.add(new BasicTsKvEntry(ts, new LongDataEntry("freeDiscSpace", getFreeDiscSpace()))); - private long getMemoryUsage() { - MemoryMXBean memoryMXBean = ManagementFactory.getMemoryMXBean(); - return memoryMXBean.getHeapMemoryUsage().getUsed(); + telemetryService.saveAndNotifyInternal(TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID, tsList, CALLBACK); } - private double getCpuUsage() { - OperatingSystemMXBean osBean = ManagementFactory.getOperatingSystemMXBean(); - return osBean.getSystemLoadAverage(); + private SystemInfoData createSystemInfoData(ServiceInfo serviceInfo) { + ProtocolStringList serviceTypes = serviceInfo.getServiceTypesList(); + SystemInfoData infoData = new SystemInfoData(); + infoData.setServiceId(serviceInfo.getServiceId()); + infoData.setServiceType(serviceTypes.size() > 1 ? "MONOLITH" : serviceTypes.get(0)); + infoData.setMemUsage(serviceInfo.getSystemInfo().getMemoryUsage()); + infoData.setCpuUsage(serviceInfo.getSystemInfo().getCpuUsage()); + infoData.setFreeDiscSpace(serviceInfo.getSystemInfo().getFreeDiscSpace()); + return infoData; } - private long getFreeDiscSpace() { - File file = new File("/"); - return file.getFreeSpace(); - } } diff --git a/common/cluster-api/src/main/proto/queue.proto b/common/cluster-api/src/main/proto/queue.proto index c1e6f4d1eb..24d1f67db3 100644 --- a/common/cluster-api/src/main/proto/queue.proto +++ b/common/cluster-api/src/main/proto/queue.proto @@ -27,6 +27,13 @@ message ServiceInfo { string serviceId = 1; repeated string serviceTypes = 2; repeated string transports = 6; + SystemInfoProto systemInfo = 10; +} + +message SystemInfoProto { + int64 memoryUsage = 1; + double cpuUsage = 2; + int64 freeDiscSpace = 3; } /** diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/SystemInfo.java b/common/data/src/main/java/org/thingsboard/server/common/data/SystemInfo.java index f1aa8b4abe..ea46c91239 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/SystemInfo.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/SystemInfo.java @@ -18,18 +18,12 @@ package org.thingsboard.server.common.data; import io.swagger.annotations.ApiModelProperty; import lombok.Data; -import java.util.Map; +import java.util.List; @Data public class SystemInfo { @ApiModelProperty(position = 1, value = "Is monolith.") private boolean isMonolith; - @ApiModelProperty(position = 2, value = "CPU usage.") - private Double cpuUsage; - @ApiModelProperty(position = 3, value = "Memory usage.") - private Long memUsage; - @ApiModelProperty(position = 4, value = "Free disc space.") - private Long freeDiscSpace; - @ApiModelProperty(position = 5, value = "Json object with info about services.") - private Map serviceInfos; + @ApiModelProperty(position = 2, value = "System data.") + private List systemData; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/SystemInfoData.java b/common/data/src/main/java/org/thingsboard/server/common/data/SystemInfoData.java new file mode 100644 index 0000000000..9345fb4feb --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/SystemInfoData.java @@ -0,0 +1,35 @@ +/** + * Copyright © 2016-2023 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data; + +import io.swagger.annotations.ApiModelProperty; +import lombok.Data; + +import java.util.Map; + +@Data +public class SystemInfoData { + @ApiModelProperty(position = 1, value = "Service Id.") + private String serviceId; + @ApiModelProperty(position = 2, value = "Service type.") + private String serviceType; + @ApiModelProperty(position = 3, value = "CPU usage.") + private Double cpuUsage; + @ApiModelProperty(position = 4, value = "Memory usage.") + private Long memUsage; + @ApiModelProperty(position = 5, value = "Free disc space.") + private Long freeDiscSpace; +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java index 0a9e00e409..b9377f3174 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java @@ -24,6 +24,7 @@ import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TbTransportService; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo; import org.thingsboard.server.queue.util.AfterContextReady; @@ -36,6 +37,10 @@ import java.util.Collections; import java.util.List; import java.util.stream.Collectors; +import static org.thingsboard.common.util.SystemUtil.getCpuUsage; +import static org.thingsboard.common.util.SystemUtil.getFreeDiscSpace; +import static org.thingsboard.common.util.SystemUtil.getMemoryUsage; + @Component @Slf4j public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider { @@ -69,11 +74,8 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider { } else { serviceTypes = Collections.singletonList(ServiceType.of(serviceType)); } - ServiceInfo.Builder builder = ServiceInfo.newBuilder() - .setServiceId(serviceId) - .addAllServiceTypes(serviceTypes.stream().map(ServiceType::name).collect(Collectors.toList())); - serviceInfo = builder.build(); + serviceInfo = getServiceInfoWithCurrentSystemInfo(); } @AfterContextReady @@ -99,4 +101,33 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider { return serviceTypes.contains(serviceType); } + @Override + public ServiceInfo getServiceInfoWithCurrentSystemInfo() { + ServiceInfo.Builder builder = ServiceInfo.newBuilder() + .setServiceId(serviceId) + .addAllServiceTypes(serviceTypes.stream().map(ServiceType::name).collect(Collectors.toList())) + .setSystemInfo(getCurrentSystemInfoProto()); + + return builder.build(); + } + + private TransportProtos.SystemInfoProto getCurrentSystemInfoProto() { + TransportProtos.SystemInfoProto.Builder builder = TransportProtos.SystemInfoProto.newBuilder(); + + Long memoryUsage = getMemoryUsage(); + if (memoryUsage != null) { + builder.setMemoryUsage(memoryUsage); + } + Double cpuUsage = getCpuUsage(); + if (cpuUsage != null) { + builder.setCpuUsage(cpuUsage); + } + Long freeDiscSpace = getFreeDiscSpace(); + if (freeDiscSpace != null) { + builder.setFreeDiscSpace(freeDiscSpace); + } + + return builder.build(); + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java index 4bc65c5e1e..dfd6e7dba5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java @@ -15,6 +15,12 @@ */ package org.thingsboard.server.queue.discovery; +import org.thingsboard.server.gen.transport.TransportProtos; + +import java.util.List; + public interface DiscoveryService { + List getOtherServers(); + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java index 7b04f3d0f6..b26284b1c4 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java @@ -22,9 +22,11 @@ import org.springframework.context.annotation.DependsOn; import org.springframework.context.event.EventListener; import org.springframework.core.annotation.Order; import org.springframework.stereotype.Service; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.AfterStartUp; import java.util.Collections; +import java.util.List; @Service @ConditionalOnProperty(prefix = "zk", value = "enabled", havingValue = "false", matchIfMissing = true) @@ -46,4 +48,8 @@ public class DummyDiscoveryService implements DiscoveryService { partitionService.recalculatePartitions(serviceInfoProvider.getServiceInfo(), Collections.emptyList()); } + @Override + public List getOtherServers() { + return Collections.emptyList(); + } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index 6888b9ad8c..776d42a334 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java @@ -17,13 +17,11 @@ package org.thingsboard.server.queue.discovery; import com.google.common.hash.HashFunction; import com.google.common.hash.Hashing; -import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.QueueId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -77,7 +75,6 @@ public class HashPartitionService implements PartitionService { private final ConcurrentMap tenantRoutingInfoMap = new ConcurrentHashMap<>(); private Map> tbTransportServicesByType = new HashMap<>(); - @Getter private List currentOtherServices; private HashFunction hashFunction; @@ -219,7 +216,7 @@ public class HashPartitionService implements PartitionService { } queueServicesMap.values().forEach(list -> list.sort(Comparator.comparing(ServiceInfo::getServiceId))); - final ConcurrentMap> newPartitions = new ConcurrentHashMap<>(); + final ConcurrentMap> newPartitions = new ConcurrentHashMap<>(); partitionSizesMap.forEach((queueKey, size) -> { for (int i = 0; i < size; i++) { ServiceInfo serviceInfo = resolveByPartitionIdx(queueServicesMap.get(queueKey), queueKey, i); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java index 3806c59a68..50e6fdf791 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java @@ -16,7 +16,6 @@ package org.thingsboard.server.queue.discovery; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.QueueId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -64,5 +63,4 @@ public interface PartitionService { void removeQueue(TransportProtos.QueueDeleteMsg queueDeleteMsg); - List getCurrentOtherServices(); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java index 4b03056253..300bc05f0a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java @@ -16,6 +16,7 @@ package org.thingsboard.server.queue.discovery; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo; public interface TbServiceInfoProvider { @@ -28,4 +29,6 @@ public interface TbServiceInfoProvider { boolean isService(ServiceType serviceType); + ServiceInfo getServiceInfoWithCurrentSystemInfo(); + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java index 63f91935d3..0a72f4064e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java @@ -16,6 +16,7 @@ package org.thingsboard.server.queue.discovery; import com.google.protobuf.InvalidProtocolBufferException; +import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.CuratorFrameworkFactory; @@ -33,21 +34,19 @@ import org.apache.zookeeper.KeeperException; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.event.ApplicationReadyEvent; -import org.springframework.context.event.EventListener; -import org.springframework.core.annotation.Order; import org.springframework.stereotype.Service; import org.springframework.util.Assert; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.discovery.event.ServiceListChangedEvent; import org.thingsboard.server.queue.util.AfterStartUp; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.List; import java.util.NoSuchElementException; -import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import static org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent.Type.CHILD_REMOVED; @@ -71,7 +70,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi private final TbServiceInfoProvider serviceInfoProvider; private final PartitionService partitionService; - private ExecutorService reconnectExecutorService; + private ScheduledExecutorService zkExecutorService; private CuratorFramework client; private PathChildrenCache cache; private String nodePath; @@ -93,7 +92,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi Assert.notNull(zkConnectionTimeout, missingProperty("zk.connection_timeout_ms")); Assert.notNull(zkSessionTimeout, missingProperty("zk.session_timeout_ms")); - reconnectExecutorService = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("zk-discovery")); + zkExecutorService = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("zk-discovery")); log.info("Initializing discovery service using ZK connect string: {}", zkUrl); @@ -101,7 +100,8 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi initZkClient(); } - private List getOtherServers() { + @Override + public List getOtherServers() { return cache.getCurrentData().stream() .filter(cd -> !cd.getPath().equals(nodePath)) .map(cd -> { @@ -128,15 +128,17 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi return; } log.info("Going to publish current server..."); - publishCurrentServer(); + zkExecutorService.scheduleAtFixedRate(this::publishCurrentServer, 0, 1, TimeUnit.MINUTES); log.info("Going to recalculate partitions..."); recalculatePartitions(); } + @SneakyThrows public synchronized void publishCurrentServer() { TransportProtos.ServiceInfo self = serviceInfoProvider.getServiceInfo(); if (currentServerExists()) { - log.info("[{}] ZK node for current instance already exists, NOT created new one: {}", self.getServiceId(), nodePath); + log.trace("[{}] Updating ZK node for current instance: {}", self.getServiceId(), nodePath); + client.setData().forPath(nodePath, serviceInfoProvider.getServiceInfoWithCurrentSystemInfo().toByteArray()); } else { try { log.info("[{}] Creating ZK node for current instance", self.getServiceId()); @@ -175,7 +177,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi return (client, newState) -> { log.info("[{}] ZK state changed: {}", self.getServiceId(), newState); if (newState == ConnectionState.LOST) { - reconnectExecutorService.submit(this::reconnect); + zkExecutorService.submit(this::reconnect); } }; } @@ -240,7 +242,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi @PreDestroy public void destroy() { destroyZkClient(); - reconnectExecutorService.shutdownNow(); + zkExecutorService.shutdownNow(); log.info("Stopped discovery service"); } @@ -280,13 +282,14 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi log.error("Failed to decode server instance for node {}", data.getPath(), e); throw e; } - log.info("Processing [{}] event for [{}]", pathChildrenCacheEvent.getType(), instance.getServiceId()); + log.debug("Processing [{}] event for [{}]", pathChildrenCacheEvent.getType(), instance.getServiceId()); switch (pathChildrenCacheEvent.getType()) { case CHILD_ADDED: - case CHILD_UPDATED: case CHILD_REMOVED: recalculatePartitions(); break; + case CHILD_UPDATED: + break; default: break; } diff --git a/common/util/src/main/java/org/thingsboard/common/util/SystemUtil.java b/common/util/src/main/java/org/thingsboard/common/util/SystemUtil.java new file mode 100644 index 0000000000..4533bbf65e --- /dev/null +++ b/common/util/src/main/java/org/thingsboard/common/util/SystemUtil.java @@ -0,0 +1,59 @@ +/** + * Copyright © 2016-2023 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.common.util; + +import lombok.extern.slf4j.Slf4j; + +import java.lang.management.ManagementFactory; +import java.lang.management.MemoryMXBean; +import java.lang.management.OperatingSystemMXBean; +import java.nio.file.FileStore; +import java.nio.file.Files; +import java.nio.file.Paths; + +@Slf4j +public class SystemUtil { + + public static Long getMemoryUsage() { + try { + MemoryMXBean memoryMXBean = ManagementFactory.getMemoryMXBean(); + return memoryMXBean.getHeapMemoryUsage().getUsed(); + } catch (Exception e) { + log.debug("Failed to get memory usage!!!", e); + } + return null; + } + + public static Double getCpuUsage() { + try { + OperatingSystemMXBean osBean = ManagementFactory.getOperatingSystemMXBean(); + return osBean.getSystemLoadAverage(); + } catch (Exception e) { + log.debug("Failed to get cpu usage!!!", e); + } + return null; + } + + public static Long getFreeDiscSpace() { + try { + FileStore store = Files.getFileStore(Paths.get("/")); + return store.getUsableSpace(); + } catch (Exception e) { + log.debug("Failed to get free disc space!!!", e); + } + return null; + } +}