From 6997cc971896a8756dd499374f0111d867d67357 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Mon, 16 Jan 2023 12:14:26 +0200 Subject: [PATCH] Monitoring service refactoring --- .../ThingsboardMonitoringApplication.java | 7 +- .../monitoring/client/WsClient.java | 2 +- .../thingsboard/monitoring/data/Latency.java | 7 ++ .../service/MonitoringReporter.java | 8 +- .../service/TransportMonitoringService.java | 75 +++++++------------ .../src/main/resources/tb-monitoring.yml | 2 +- 6 files changed, 44 insertions(+), 57 deletions(-) diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java b/monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java index cbfba4da51..aa9d24e319 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java @@ -103,13 +103,8 @@ public class ThingsboardMonitoringApplication { } @Bean - public ScheduledExecutorService monitoringExecutor(@Value("${monitoring.monitoring_thread_pool_size}") int threadPoolSize) { + public ScheduledExecutorService monitoringExecutor(@Value("${monitoring.monitoring_executor_thread_pool_size}") int threadPoolSize) { return Executors.newScheduledThreadPool(threadPoolSize, ThingsBoardThreadFactory.forName("monitoring-executor")); } - @Bean - public ExecutorService requestExecutor() { - return Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("request-executor")); - } - } diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java b/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java index 56d29de5fa..ea1a5035f2 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java @@ -37,7 +37,7 @@ import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @Slf4j -public class WsClient extends WebSocketClient { +public class WsClient extends WebSocketClient implements AutoCloseable { public volatile String lastMsg; private CountDownLatch reply; diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/data/Latency.java b/monitoring/src/main/java/org/thingsboard/monitoring/data/Latency.java index 90914f09e1..6c385cf6e4 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/data/Latency.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/data/Latency.java @@ -49,6 +49,13 @@ public class Latency { return key; } + public synchronized Latency snapshot() { + Latency snapshot = new Latency(key); + snapshot.latencySum.set(latencySum.get()); + snapshot.counter.set(counter.get()); + return snapshot; + } + @Override public String toString() { return "Latency{" + diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java b/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java index 8ca3883110..5721609d27 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java @@ -70,9 +70,14 @@ public class MonitoringReporter { @EventListener(ApplicationReadyEvent.class) public void startLatenciesMonitoring() { - monitoringExecutor.scheduleWithFixedDelay(() -> { + monitoringExecutor.scheduleAtFixedRate(() -> { List latencies = this.latencies.values().stream() .filter(Latency::isNotEmpty) + .map(latency -> { + Latency snapshot = latency.snapshot(); + latency.reset(); + return snapshot; + }) .collect(Collectors.toList()); if (latencies.isEmpty()) { return; @@ -95,7 +100,6 @@ public class MonitoringReporter { ObjectNode msg = JacksonUtil.newObjectNode(); latencies.forEach(latency -> { msg.set(latency.getKey(), new DoubleNode(latency.getAvg())); - latency.reset(); }); tbClient.saveEntityTelemetry(entityId, "time", msg); } catch (Exception e) { diff --git a/monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java b/monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java index 4d0455628d..4fe003d87d 100644 --- a/monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java +++ b/monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java @@ -36,12 +36,8 @@ import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.Optional; import java.util.UUID; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Future; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; @Slf4j public abstract class TransportMonitoringService { @@ -59,8 +55,6 @@ public abstract class TransportMonitoringService { - WsClient wsClient = null; - try { - log.trace("[{}] Checking", transportInfo); - wsClient = establishWsClient(); - wsClient.registerWaitForUpdate(); - - String testValue = UUID.randomUUID().toString(); - String testPayload = JacksonUtil.newObjectNode().set(TEST_TELEMETRY_KEY, new TextNode(testValue)).toString(); - try { - initClientAndSendPayload(testPayload); - log.trace("[{}] Sent test payload ({})", transportInfo, testPayload); - } catch (Throwable e) { - throw new TransportFailureException(e); - } - - log.trace("[{}] Waiting for WS update", transportInfo); - checkWsUpdate(wsClient, testValue); - - monitoringReporter.serviceIsOk(transportInfo); - monitoringReporter.serviceIsOk(MonitoredServiceKey.GENERAL); - } catch (TransportFailureException transportFailureException) { - monitoringReporter.serviceFailure(transportInfo, transportFailureException); - } catch (Exception e) { - monitoringReporter.serviceFailure(MonitoredServiceKey.GENERAL, e); - } finally { - if (wsClient != null) wsClient.close(); - } + check(); }, config.getInitialDelayMs(), config.getMonitoringRateMs(), TimeUnit.MILLISECONDS); log.info("Started monitoring for transport type {} for target {}", getTransportType(), target); } - private void initClientAndSendPayload(String payload) throws Throwable { - initClient(); - stopWatch.start(); - Future resultFuture = requestExecutor.submit(() -> { + private void check() { + log.trace("[{}] Checking", transportInfo); + try (WsClient wsClient = establishWsClient()) { + wsClient.registerWaitForUpdate(); + + String testValue = UUID.randomUUID().toString(); + String testPayload = JacksonUtil.newObjectNode().set(TEST_TELEMETRY_KEY, new TextNode(testValue)).toString(); try { - sendTestPayload(payload); - } catch (Exception e) { - throw new RuntimeException(e.getMessage(), e); + initClientAndSendPayload(testPayload); + log.trace("[{}] Sent test payload ({})", transportInfo, testPayload); + } catch (Throwable e) { + throw new TransportFailureException(e); } - }); - try { - resultFuture.get(config.getRequestTimeoutMs(), TimeUnit.MILLISECONDS); - } catch (ExecutionException e) { - throw e.getCause(); - } catch (TimeoutException e) { - throw new TimeoutException("Transport request timeout"); + + log.trace("[{}] Waiting for WS update", transportInfo); + checkWsUpdate(wsClient, testValue); + + monitoringReporter.serviceIsOk(transportInfo); + monitoringReporter.serviceIsOk(MonitoredServiceKey.GENERAL); + } catch (TransportFailureException transportFailureException) { + monitoringReporter.serviceFailure(transportInfo, transportFailureException); + } catch (Exception e) { + monitoringReporter.serviceFailure(MonitoredServiceKey.GENERAL, e); } + } + + private void initClientAndSendPayload(String payload) throws Throwable { + initClient(); + stopWatch.start(); + sendTestPayload(payload); monitoringReporter.reportLatency(Latencies.transportRequest(getTransportType()), stopWatch.getTime()); } diff --git a/monitoring/src/main/resources/tb-monitoring.yml b/monitoring/src/main/resources/tb-monitoring.yml index 890c2ab10d..3d71a7dbfa 100644 --- a/monitoring/src/main/resources/tb-monitoring.yml +++ b/monitoring/src/main/resources/tb-monitoring.yml @@ -74,4 +74,4 @@ monitoring: reporting_entity_type: '${LATENCY_REPORTING_ENTITY_TYPE:ASSET}' reporting_entity_id: '${LATENCY_REPORTING_ENTITY_ID:}' - monitoring_thread_pool_size: '${MONITORING_THREAD_POOL_SIZE:10}' + monitoring_executor_thread_pool_size: '${MONITORING_EXECUTOR_THREAD_POOL_SIZE:10}'