Browse Source

Monitoring service refactoring

pull/7725/head
ViacheslavKlimov 4 years ago
parent
commit
6997cc9718
  1. 7
      monitoring/src/main/java/org/thingsboard/monitoring/ThingsboardMonitoringApplication.java
  2. 2
      monitoring/src/main/java/org/thingsboard/monitoring/client/WsClient.java
  3. 7
      monitoring/src/main/java/org/thingsboard/monitoring/data/Latency.java
  4. 8
      monitoring/src/main/java/org/thingsboard/monitoring/service/MonitoringReporter.java
  5. 75
      monitoring/src/main/java/org/thingsboard/monitoring/service/TransportMonitoringService.java
  6. 2
      monitoring/src/main/resources/tb-monitoring.yml

7
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"));
}
}

2
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;

7
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{" +

8
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<Latency> 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) {

75
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<C extends TransportMonitoringServiceConfig> {
@ -59,8 +55,6 @@ public abstract class TransportMonitoringService<C extends TransportMonitoringSe
@Autowired
private ScheduledExecutorService monitoringExecutor;
@Autowired
private ExecutorService requestExecutor;
@Autowired
private TbClient tbClient;
@Autowired
private TbStopWatch stopWatch;
@ -79,54 +73,41 @@ public abstract class TransportMonitoringService<C extends TransportMonitoringSe
public final void startMonitoring() {
monitoringExecutor.scheduleWithFixedDelay(() -> {
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());
}

2
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}'

Loading…
Cancel
Save