diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index b20e53893b..fe190ec234 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -174,17 +174,23 @@ public class DefaultTbQueueRequestTemplate responseMetaData = new ResponseMetaData<>(tickTs + maxRequestTimeout, future); pendingRequests.putIfAbsent(requestId, responseMetaData); log.trace("[{}] Sending request, key [{}], expTime [{}]", requestId, request.getKey(), responseMetaData.expTime); - if (messagesStats != null) messagesStats.incrementTotal(); + if (messagesStats != null) { + messagesStats.incrementTotal(); + } requestTemplate.send(TopicPartitionInfo.builder().topic(requestTemplate.getDefaultTopic()).build(), request, new TbQueueCallback() { @Override public void onSuccess(TbQueueMsgMetadata metadata) { - if (messagesStats != null) messagesStats.incrementSuccessful(); + if (messagesStats != null) { + messagesStats.incrementSuccessful(); + } log.trace("[{}] Request sent: {}", requestId, metadata); } @Override public void onFailure(Throwable t) { - if (messagesStats != null) messagesStats.incrementFailed(); + if (messagesStats != null) { + messagesStats.incrementFailed(); + } pendingRequests.remove(requestId); future.setException(t); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java index 327b8d675b..22fda66759 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java @@ -38,10 +38,11 @@ public class TbSqlBlockingQueue implements TbSqlQueue { private final TbSqlBlockingQueueParams params; private ExecutorService executor; - private MessagesStats stats; + private final MessagesStats stats; - public TbSqlBlockingQueue(TbSqlBlockingQueueParams params) { + public TbSqlBlockingQueue(TbSqlBlockingQueueParams params, MessagesStats stats) { this.params = params; + this.stats = stats; } @Override @@ -108,7 +109,7 @@ public class TbSqlBlockingQueue implements TbSqlQueue { public ListenableFuture add(E element) { SettableFuture future = SettableFuture.create(); queue.add(new TbSqlQueueElement<>(future, element)); - params.getStats().incrementTotal(); + stats.incrementTotal(); return future; } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueParams.java b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueParams.java index 330e60f336..a63461787e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueParams.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueParams.java @@ -19,6 +19,7 @@ import lombok.Builder; import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.stats.MessagesStats; +import org.thingsboard.server.common.stats.StatsFactory; @Slf4j @Data @@ -29,5 +30,5 @@ public class TbSqlBlockingQueueParams { private final int batchSize; private final long maxDelay; private final long statsPrintIntervalMs; - private final MessagesStats stats; + private final String statsNamePrefix; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java index 2c53cc4ad9..d9596a6efd 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java @@ -18,6 +18,8 @@ package org.thingsboard.server.dao.sql; import com.google.common.util.concurrent.ListenableFuture; import lombok.Data; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.stats.MessagesStats; +import org.thingsboard.server.common.stats.StatsFactory; import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; @@ -32,10 +34,12 @@ public class TbSqlBlockingQueueWrapper { private ScheduledLogExecutorComponent logExecutor; private final Function hashCodeFunction; private final int maxThreads; + private final StatsFactory statsFactory; public void init(ScheduledLogExecutorComponent logExecutor, Consumer> saveFunction) { for (int i = 0; i < maxThreads; i++) { - TbSqlBlockingQueue queue = new TbSqlBlockingQueue<>(params); + MessagesStats stats = statsFactory.createMessagesStats(params.getStatsNamePrefix() + ".queue." + i); + TbSqlBlockingQueue queue = new TbSqlBlockingQueue<>(params, stats); queues.add(queue); queue.init(logExecutor, saveFunction, i); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java index 5bd5d4876f..2a38c6a0f1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java @@ -85,11 +85,11 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl .batchSize(batchSize) .maxDelay(maxDelay) .statsPrintIntervalMs(statsPrintIntervalMs) - .stats(statsFactory.createMessagesStats("attributes")) + .statsNamePrefix("attributes") .build(); Function hashcodeFunction = entity -> entity.getId().getEntityId().hashCode(); - queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads); + queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory); queue.init(logExecutor, v -> attributeKvInsertRepository.saveOrUpdate(v)); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java index 39c958a45c..459a0bc015 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java @@ -46,6 +46,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; +import java.util.function.Function; import java.util.stream.Collectors; @Slf4j @@ -69,11 +70,11 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq .batchSize(tsBatchSize) .maxDelay(tsMaxDelay) .statsPrintIntervalMs(tsStatsPrintIntervalMs) - .stats(statsFactory.createMessagesStats("ts")) + .statsNamePrefix("ts") .build(); Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); - tsQueue = new TbSqlBlockingQueueWrapper<>(tsParams, hashcodeFunction, tsBatchThreads); + tsQueue = new TbSqlBlockingQueueWrapper<>(tsParams, hashcodeFunction, tsBatchThreads, statsFactory); tsQueue.init(logExecutor, v -> insertRepository.saveOrUpdate(v)); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java index bad2bb53e4..9710edad4a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java @@ -128,11 +128,11 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx .batchSize(tsLatestBatchSize) .maxDelay(tsLatestMaxDelay) .statsPrintIntervalMs(tsLatestStatsPrintIntervalMs) - .stats(statsFactory.createMessagesStats("ts.latest")) + .statsNamePrefix("ts.latest") .build(); java.util.function.Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); - tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads); + tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads, statsFactory); tsLatestQueue.init(logExecutor, v -> { Map trueLatest = new HashMap<>(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java index a55da5a448..e3ad6bd2bc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java @@ -78,11 +78,11 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements .batchSize(tsBatchSize) .maxDelay(tsMaxDelay) .statsPrintIntervalMs(tsStatsPrintIntervalMs) - .stats(statsFactory.createMessagesStats("ts.timescale")) + .statsNamePrefix("ts.timescale") .build(); Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); - tsQueue = new TbSqlBlockingQueueWrapper<>(tsParams, hashcodeFunction, timescaleBatchThreads); + tsQueue = new TbSqlBlockingQueueWrapper<>(tsParams, hashcodeFunction, timescaleBatchThreads, statsFactory); tsQueue.init(logExecutor, v -> insertRepository.saveOrUpdate(v)); }