From d98bbacd43de4d0d9b08f751bd638be94f1859ac Mon Sep 17 00:00:00 2001 From: vzikratyi Date: Fri, 10 Jul 2020 13:02:46 +0300 Subject: [PATCH] Added stats to TbSqlQueue and BufferedRateExecutor --- .../common/msg/stats/DefaultCounter.java | 48 ++++++++++ .../msg/stats/DefaultMessagesStats.java | 22 +++++ .../common/msg/stats/DefaultStatsFactory.java | 16 ++++ .../common/msg/stats/MessagesStats.java | 9 +- .../server/common/msg/stats/StatsCounter.java | 25 +---- .../server/common/msg/stats/StatsFactory.java | 4 + .../server/common/msg/stats/StatsType.java | 2 +- .../nosql/CassandraBufferedRateExecutor.java | 77 +++++++++------- .../server/dao/sql/TbSqlBlockingQueue.java | 16 ++-- .../dao/sql/TbSqlBlockingQueueParams.java | 2 + .../dao/sql/attributes/JpaAttributeDao.java | 5 + ...stractChunkedAggregationTimeseriesDao.java | 6 +- .../dao/sqlts/AbstractSqlTimeseriesDao.java | 5 + .../timescale/TimescaleTimeseriesDao.java | 5 + .../util/AbstractBufferedRateExecutor.java | 48 +++++----- .../dao/util/BufferedRateExecutorStats.java | 91 +++++++++++++++++++ 16 files changed, 294 insertions(+), 87 deletions(-) create mode 100644 common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultCounter.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateExecutorStats.java diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultCounter.java b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultCounter.java new file mode 100644 index 0000000000..b6fe419c51 --- /dev/null +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultCounter.java @@ -0,0 +1,48 @@ +/** + * Copyright © 2016-2020 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.msg.stats; + +import io.micrometer.core.instrument.Counter; + +import java.util.concurrent.atomic.AtomicInteger; + +public class DefaultCounter { + private final AtomicInteger aiCounter; + private final Counter micrometerCounter; + + public DefaultCounter(AtomicInteger aiCounter, Counter micrometerCounter) { + this.aiCounter = aiCounter; + this.micrometerCounter = micrometerCounter; + } + + public void increment() { + aiCounter.incrementAndGet(); + micrometerCounter.increment(); + } + + public void clear() { + aiCounter.set(0); + } + + public int get() { + return aiCounter.get(); + } + + public void add(int delta){ + aiCounter.addAndGet(delta); + micrometerCounter.increment(delta); + } +} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultMessagesStats.java b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultMessagesStats.java index 29990eaea4..12df2593b6 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultMessagesStats.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultMessagesStats.java @@ -40,4 +40,26 @@ public class DefaultMessagesStats implements MessagesStats { public void incrementFailed(int amount) { failedCounter.add(amount); } + + @Override + public int getTotal() { + return totalCounter.get(); + } + + @Override + public int getSuccessful() { + return successfulCounter.get(); + } + + @Override + public int getFailed() { + return failedCounter.get(); + } + + @Override + public void reset() { + totalCounter.clear(); + successfulCounter.clear(); + failedCounter.clear(); + } } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultStatsFactory.java b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultStatsFactory.java index e34f6e0ff2..3b084729cb 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultStatsFactory.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultStatsFactory.java @@ -17,6 +17,7 @@ package org.thingsboard.server.common.msg.stats; import io.micrometer.core.instrument.Counter; import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Tags; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; @@ -50,6 +51,21 @@ public class DefaultStatsFactory implements StatsFactory { ); } + @Override + public DefaultCounter createDefaultCounter(String key, String... tags) { + return new DefaultCounter( + new AtomicInteger(0), + metricsEnabled ? + meterRegistry.counter(key, tags) + : STUB_COUNTER + ); + } + + @Override + public T createGauge(String key, T number, String... tags) { + return meterRegistry.gauge(key, Tags.of(tags), number); + } + @Override public MessagesStats createMessagesStats(String key) { StatsCounter totalCounter = createStatsCounter(key, TOTAL_MSGS); diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/MessagesStats.java b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/MessagesStats.java index 077d2ce01a..cc080c5dd4 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/MessagesStats.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/MessagesStats.java @@ -28,10 +28,17 @@ public interface MessagesStats { void incrementSuccessful(int amount); - default void incrementFailed() { incrementFailed(1); } void incrementFailed(int amount); + + int getTotal(); + + int getSuccessful(); + + int getFailed(); + + void reset(); } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsCounter.java b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsCounter.java index 0097949cb7..971da18e1e 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsCounter.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsCounter.java @@ -19,35 +19,14 @@ import io.micrometer.core.instrument.Counter; import java.util.concurrent.atomic.AtomicInteger; -public class StatsCounter { - private final AtomicInteger aiCounter; - private final Counter micrometerCounter; +public class StatsCounter extends DefaultCounter { private final String name; public StatsCounter(AtomicInteger aiCounter, Counter micrometerCounter, String name) { - this.aiCounter = aiCounter; - this.micrometerCounter = micrometerCounter; + super(aiCounter, micrometerCounter); this.name = name; } - public void increment() { - aiCounter.incrementAndGet(); - micrometerCounter.increment(); - } - - public void clear() { - aiCounter.set(0); - } - - public int get() { - return aiCounter.get(); - } - - public void add(int delta){ - aiCounter.addAndGet(delta); - micrometerCounter.increment(delta); - } - public String getName() { return name; } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsFactory.java b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsFactory.java index 243c1e05a0..b21ea917ce 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsFactory.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsFactory.java @@ -18,5 +18,9 @@ package org.thingsboard.server.common.msg.stats; public interface StatsFactory { StatsCounter createStatsCounter(String key, String statsName); + DefaultCounter createDefaultCounter(String key, String... tags); + + T createGauge(String key, T number, String... tags); + MessagesStats createMessagesStats(String key); } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsType.java b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsType.java index 825f9fcac8..c5a6c1f3d2 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsType.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsType.java @@ -16,7 +16,7 @@ package org.thingsboard.server.common.msg.stats; public enum StatsType { - RULE_ENGINE("ruleEngine"), CORE("core"), TRANSPORT("transport"), JS_INVOKE("jsInvoke"); + RULE_ENGINE("ruleEngine"), CORE("core"), TRANSPORT("transport"), JS_INVOKE("jsInvoke"), RATE_EXECUTOR("rateExecutor"); private String name; diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateExecutor.java index 99e1b9ccc5..9b55cf0b30 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateExecutor.java @@ -24,6 +24,9 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.stats.DefaultCounter; +import org.thingsboard.server.common.msg.stats.StatsCounter; +import org.thingsboard.server.common.msg.stats.StatsFactory; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.util.AbstractBufferedRateExecutor; import org.thingsboard.server.dao.util.AsyncTaskContext; @@ -57,48 +60,58 @@ public class CassandraBufferedRateExecutor extends AbstractBufferedRateExecutor< @Value("${cassandra.query.tenant_rate_limits.enabled}") boolean tenantRateLimitsEnabled, @Value("${cassandra.query.tenant_rate_limits.configuration}") String tenantRateLimitsConfiguration, @Value("${cassandra.query.tenant_rate_limits.print_tenant_names}") boolean printTenantNames, - @Value("${cassandra.query.print_queries_freq:0}") int printQueriesFreq) { - super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, tenantRateLimitsEnabled, tenantRateLimitsConfiguration, printQueriesFreq); + @Value("${cassandra.query.print_queries_freq:0}") int printQueriesFreq, + @Autowired StatsFactory statsFactory) { + super(queueLimit, concurrencyLimit, maxWaitTime, dispatcherThreads, callbackThreads, pollMs, tenantRateLimitsEnabled, tenantRateLimitsConfiguration, printQueriesFreq, statsFactory); this.printTenantNames = printTenantNames; } @Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") public void printStats() { int queueSize = getQueueSize(); - int totalAddedValue = totalAdded.getAndSet(0); - int totalLaunchedValue = totalLaunched.getAndSet(0); - int totalReleasedValue = totalReleased.getAndSet(0); - int totalFailedValue = totalFailed.getAndSet(0); - int totalExpiredValue = totalExpired.getAndSet(0); - int totalRejectedValue = totalRejected.getAndSet(0); - int totalRateLimitedValue = totalRateLimited.getAndSet(0); - int rateLimitedTenantsValue = rateLimitedTenants.size(); - int concurrencyLevelValue = concurrencyLevel.get(); - if (queueSize > 0 || totalAddedValue > 0 || totalLaunchedValue > 0 || totalReleasedValue > 0 || - totalFailedValue > 0 || totalExpiredValue > 0 || totalRejectedValue > 0 || totalRateLimitedValue > 0 || rateLimitedTenantsValue > 0 - || concurrencyLevelValue > 0) { - log.info("Permits queueSize [{}] totalAdded [{}] totalLaunched [{}] totalReleased [{}] totalFailed [{}] totalExpired [{}] totalRejected [{}] " + - "totalRateLimited [{}] totalRateLimitedTenants [{}] currBuffer [{}] ", - queueSize, totalAddedValue, totalLaunchedValue, totalReleasedValue, - totalFailedValue, totalExpiredValue, totalRejectedValue, totalRateLimitedValue, rateLimitedTenantsValue, concurrencyLevelValue); + int rateLimitedTenantsCount = (int) stats.getRateLimitedTenants().values().stream() + .filter(defaultCounter -> defaultCounter.get() > 0) + .count(); + + if (queueSize > 0 + || rateLimitedTenantsCount > 0 + || concurrencyLevel.get() > 0 + || stats.getStatsCounters().stream().anyMatch(counter -> counter.get() > 0) + ) { + StringBuilder statsBuilder = new StringBuilder(); + + statsBuilder.append("queueSize").append(" = [").append(queueSize).append("] "); + stats.getStatsCounters().forEach(counter -> { + statsBuilder.append(counter.getName()).append(" = [").append(counter.get()).append("] "); + }); + statsBuilder.append("totalRateLimitedTenants").append(" = [").append(rateLimitedTenantsCount).append("] "); + statsBuilder.append(CONCURRENCY_LEVEL).append(" = [").append(concurrencyLevel.get()).append("] "); + + stats.getStatsCounters().forEach(StatsCounter::clear); + log.info("Permits {}", statsBuilder); } - rateLimitedTenants.forEach(((tenantId, counter) -> { - if (printTenantNames) { - String name = tenantNamesCache.computeIfAbsent(tenantId, tId -> { - try { - return entityService.fetchEntityNameAsync(TenantId.SYS_TENANT_ID, tenantId).get(); - } catch (Exception e) { - log.error("[{}] Failed to get tenant name", tenantId, e); - return "N/A"; + stats.getRateLimitedTenants().entrySet().stream() + .filter(entry -> entry.getValue().get() > 0) + .forEach(entry -> { + TenantId tenantId = entry.getKey(); + DefaultCounter counter = entry.getValue(); + int rateLimitedRequests = counter.get(); + counter.clear(); + if (printTenantNames) { + String name = tenantNamesCache.computeIfAbsent(tenantId, tId -> { + try { + return entityService.fetchEntityNameAsync(TenantId.SYS_TENANT_ID, tenantId).get(); + } catch (Exception e) { + log.error("[{}] Failed to get tenant name", tenantId, e); + return "N/A"; + } + }); + log.info("[{}][{}] Rate limited requests: {}", tenantId, name, rateLimitedRequests); + } else { + log.info("[{}] Rate limited requests: {}", tenantId, rateLimitedRequests); } }); - log.info("[{}][{}] Rate limited requests: {}", tenantId, name, counter); - } else { - log.info("[{}] Rate limited requests: {}", tenantId, counter); - } - })); - rateLimitedTenants.clear(); } @PreDestroy 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 3554fe7bce..69ee19e586 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 @@ -19,6 +19,7 @@ import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.common.msg.stats.MessagesStats; import java.util.ArrayList; import java.util.List; @@ -35,12 +36,10 @@ import java.util.stream.Collectors; public class TbSqlBlockingQueue implements TbSqlQueue { private final BlockingQueue> queue = new LinkedBlockingQueue<>(); - private final AtomicInteger addedCount = new AtomicInteger(); - private final AtomicInteger savedCount = new AtomicInteger(); - private final AtomicInteger failedCount = new AtomicInteger(); private final TbSqlBlockingQueueParams params; private ExecutorService executor; + private MessagesStats stats; public TbSqlBlockingQueue(TbSqlBlockingQueueParams params) { this.params = params; @@ -68,7 +67,7 @@ public class TbSqlBlockingQueue implements TbSqlQueue { log.debug("[{}] Going to save {} entities", logName, entities.size()); saveFunction.accept(entities.stream().map(TbSqlQueueElement::getEntity).collect(Collectors.toList())); entities.forEach(v -> v.getFuture().set(null)); - savedCount.addAndGet(entities.size()); + stats.incrementSuccessful(entities.size()); if (!fullPack) { long remainingDelay = maxDelay - (System.currentTimeMillis() - currentTs); if (remainingDelay > 0) { @@ -76,7 +75,7 @@ public class TbSqlBlockingQueue implements TbSqlQueue { } } } catch (Exception e) { - failedCount.addAndGet(entities.size()); + stats.incrementFailed(entities.size()); entities.forEach(entityFutureWrapper -> entityFutureWrapper.getFuture().setException(e)); if (e instanceof InterruptedException) { log.info("[{}] Queue polling was interrupted", logName); @@ -91,9 +90,10 @@ public class TbSqlBlockingQueue implements TbSqlQueue { }); logExecutor.scheduleAtFixedRate(() -> { - if (queue.size() > 0 || addedCount.get() > 0 || savedCount.get() > 0 || failedCount.get() > 0) { + if (queue.size() > 0 || stats.getTotal() > 0 || stats.getSuccessful() > 0 || stats.getFailed() > 0) { log.info("Queue-{} [{}] queueSize [{}] totalAdded [{}] totalSaved [{}] totalFailed [{}]", index, - params.getLogName(), queue.size(), addedCount.getAndSet(0), savedCount.getAndSet(0), failedCount.getAndSet(0)); + params.getLogName(), queue.size(), stats.getTotal(), stats.getSuccessful(), stats.getFailed()); + stats.reset(); } }, params.getStatsPrintIntervalMs(), params.getStatsPrintIntervalMs(), TimeUnit.MILLISECONDS); } @@ -109,7 +109,7 @@ public class TbSqlBlockingQueue implements TbSqlQueue { public ListenableFuture add(E element) { SettableFuture future = SettableFuture.create(); queue.add(new TbSqlQueueElement<>(future, element)); - addedCount.incrementAndGet(); + params.getStats().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 9e7c76ace5..8bebe6646c 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 @@ -18,6 +18,7 @@ package org.thingsboard.server.dao.sql; import lombok.Builder; import lombok.Data; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.msg.stats.MessagesStats; @Slf4j @Data @@ -28,4 +29,5 @@ public class TbSqlBlockingQueueParams { private final int batchSize; private final long maxDelay; private final long statsPrintIntervalMs; + private final MessagesStats stats; } 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 a3b17c7863..d9fedb28b4 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 @@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.UUIDConverter; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.msg.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.attributes.AttributesDao; import org.thingsboard.server.dao.model.sql.AttributeKvCompositeKey; @@ -60,6 +61,9 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl @Autowired private AttributeKvInsertRepository attributeKvInsertRepository; + @Autowired + private StatsFactory statsFactory; + @Value("${sql.attributes.batch_size:1000}") private int batchSize; @@ -81,6 +85,7 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl .batchSize(batchSize) .maxDelay(maxDelay) .statsPrintIntervalMs(statsPrintIntervalMs) + .stats(statsFactory.createMessagesStats("attributes")) .build(); Function hashcodeFunction = entity -> entity.getId().getEntityId().hashCode(); 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 84dddcaced..d1adb2452c 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 @@ -30,8 +30,10 @@ import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.msg.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity; +import org.thingsboard.server.dao.sql.TbSqlBlockingQueue; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.sqlts.insert.InsertTsRepository; @@ -44,7 +46,6 @@ 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 @@ -57,6 +58,8 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq protected InsertTsRepository insertRepository; protected TbSqlBlockingQueueWrapper tsQueue; + @Autowired + private StatsFactory statsFactory; @PostConstruct protected void init() { @@ -66,6 +69,7 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq .batchSize(tsBatchSize) .maxDelay(tsMaxDelay) .statsPrintIntervalMs(tsStatsPrintIntervalMs) + .stats(statsFactory.createMessagesStats("ts")) .build(); Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); 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 75b3d6951b..b378934fe5 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 @@ -34,6 +34,7 @@ import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.msg.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionary; import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey; @@ -102,6 +103,9 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx @Autowired protected ScheduledLogExecutorComponent logExecutor; + @Autowired + private StatsFactory statsFactory; + @Value("${sql.ts.batch_size:1000}") protected int tsBatchSize; @@ -124,6 +128,7 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx .batchSize(tsLatestBatchSize) .maxDelay(tsLatestMaxDelay) .statsPrintIntervalMs(tsLatestStatsPrintIntervalMs) + .stats(statsFactory.createMessagesStats("ts.latest")) .build(); java.util.function.Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); 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 a84a3e884f..2d43177a2c 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 @@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.msg.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.model.sqlts.timescale.ts.TimescaleTsKvEntity; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; @@ -61,6 +62,9 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements @Autowired private AggregationRepository aggregationRepository; + @Autowired + private StatsFactory statsFactory; + @Autowired protected InsertTsRepository insertRepository; @@ -74,6 +78,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements .batchSize(tsBatchSize) .maxDelay(tsMaxDelay) .statsPrintIntervalMs(tsStatsPrintIntervalMs) + .stats(statsFactory.createMessagesStats("ts.timescale")) .build(); Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java b/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java index ea8d4b6fa1..94894bd83c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java @@ -23,10 +23,16 @@ import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.stats.DefaultCounter; +import org.thingsboard.server.common.msg.stats.StatsCounter; +import org.thingsboard.server.common.msg.stats.StatsFactory; +import org.thingsboard.server.common.msg.stats.StatsType; import org.thingsboard.server.common.msg.tools.TbRateLimits; import org.thingsboard.server.dao.nosql.CassandraStatementTask; import javax.annotation.Nullable; +import java.util.ArrayList; +import java.util.List; import java.util.UUID; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; @@ -38,6 +44,8 @@ import java.util.regex.Matcher; @Slf4j public abstract class AbstractBufferedRateExecutor, V> implements BufferedRateExecutor { + public static final String CONCURRENCY_LEVEL = "currBuffer"; + private final long maxWaitTime; private final long pollMs; private final BlockingQueue> queue; @@ -49,20 +57,14 @@ public abstract class AbstractBufferedRateExecutor perTenantLimits = new ConcurrentHashMap<>(); - protected final ConcurrentMap rateLimitedTenants = new ConcurrentHashMap<>(); - - protected final AtomicInteger concurrencyLevel = new AtomicInteger(); - protected final AtomicInteger totalAdded = new AtomicInteger(); - protected final AtomicInteger totalLaunched = new AtomicInteger(); - protected final AtomicInteger totalReleased = new AtomicInteger(); - protected final AtomicInteger totalFailed = new AtomicInteger(); - protected final AtomicInteger totalExpired = new AtomicInteger(); - protected final AtomicInteger totalRejected = new AtomicInteger(); - protected final AtomicInteger totalRateLimited = new AtomicInteger(); - protected final AtomicInteger printQueriesIdx = new AtomicInteger(); + + private final AtomicInteger printQueriesIdx = new AtomicInteger(0); + + protected final AtomicInteger concurrencyLevel; + protected final BufferedRateExecutorStats stats; public AbstractBufferedRateExecutor(int queueLimit, int concurrencyLimit, long maxWaitTime, int dispatcherThreads, int callbackThreads, long pollMs, - boolean perTenantLimitsEnabled, String perTenantLimitsConfiguration, int printQueriesFreq) { + boolean perTenantLimitsEnabled, String perTenantLimitsConfiguration, int printQueriesFreq, StatsFactory statsFactory) { this.maxWaitTime = maxWaitTime; this.pollMs = pollMs; this.concurrencyLimit = concurrencyLimit; @@ -73,6 +75,10 @@ public abstract class AbstractBufferedRateExecutor new TbRateLimits(perTenantLimitsConfiguration)); if (!rateLimits.tryConsume()) { - rateLimitedTenants.computeIfAbsent(task.getTenantId(), tId -> new AtomicInteger(0)).incrementAndGet(); - totalRateLimited.incrementAndGet(); + stats.incrementRateLimitedTenant(task.getTenantId()); + stats.getTotalRateLimited().increment(); settableFuture.setException(new TenantRateLimitException()); perTenantLimitReached = true; } @@ -98,10 +104,10 @@ public abstract class AbstractBufferedRateExecutor(UUID.randomUUID(), task, settableFuture, System.currentTimeMillis())); } catch (IllegalStateException e) { - totalRejected.incrementAndGet(); + stats.getTotalRejected().increment(); settableFuture.setException(e); } } @@ -146,14 +152,14 @@ public abstract class AbstractBufferedRateExecutor 0) { - totalLaunched.incrementAndGet(); + stats.getTotalLaunched().increment(); ListenableFuture result = execute(finalTaskCtx); result = Futures.withTimeout(result, timeout, TimeUnit.MILLISECONDS, timeoutExecutor); Futures.addCallback(result, new FutureCallback() { @Override public void onSuccess(@Nullable V result) { logTask("Releasing", finalTaskCtx); - totalReleased.incrementAndGet(); + stats.getTotalReleased().increment(); concurrencyLevel.decrementAndGet(); finalTaskCtx.getFuture().set(result); } @@ -165,7 +171,7 @@ public abstract class AbstractBufferedRateExecutor rateLimitedTenants = new ConcurrentHashMap<>(); + + private final List statsCounters = new ArrayList<>(); + + private final StatsCounter totalAdded; + private final StatsCounter totalLaunched; + private final StatsCounter totalReleased; + private final StatsCounter totalFailed; + private final StatsCounter totalExpired; + private final StatsCounter totalRejected; + private final StatsCounter totalRateLimited; + + public BufferedRateExecutorStats(StatsFactory statsFactory) { + this.statsFactory = statsFactory; + + String key = StatsType.RATE_EXECUTOR.getName(); + + this.totalAdded = statsFactory.createStatsCounter(key, TOTAL_ADDED); + this.totalLaunched = statsFactory.createStatsCounter(key, TOTAL_LAUNCHED); + this.totalReleased = statsFactory.createStatsCounter(key, TOTAL_RELEASED); + this.totalFailed = statsFactory.createStatsCounter(key, TOTAL_FAILED); + this.totalExpired = statsFactory.createStatsCounter(key, TOTAL_EXPIRED); + this.totalRejected = statsFactory.createStatsCounter(key, TOTAL_REJECTED); + this.totalRateLimited = statsFactory.createStatsCounter(key, TOTAL_RATE_LIMITED); + + this.statsCounters.add(totalAdded); + this.statsCounters.add(totalLaunched); + this.statsCounters.add(totalReleased); + this.statsCounters.add(totalFailed); + this.statsCounters.add(totalExpired); + this.statsCounters.add(totalRejected); + this.statsCounters.add(totalRateLimited); + } + + public void incrementRateLimitedTenant(TenantId tenantId){ + rateLimitedTenants.computeIfAbsent(tenantId, + tId -> { + String key = StatsType.RATE_EXECUTOR.getName() + ".tenant"; + return statsFactory.createDefaultCounter(key, TENANT_ID_TAG, tId.toString()); + } + ) + .increment(); + } +}