Browse Source

Added stats to TbSqlQueue and BufferedRateExecutor

pull/3199/head
vzikratyi 6 years ago
committed by Andrew Shvayka
parent
commit
d98bbacd43
  1. 48
      common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultCounter.java
  2. 22
      common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultMessagesStats.java
  3. 16
      common/message/src/main/java/org/thingsboard/server/common/msg/stats/DefaultStatsFactory.java
  4. 9
      common/message/src/main/java/org/thingsboard/server/common/msg/stats/MessagesStats.java
  5. 25
      common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsCounter.java
  6. 4
      common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsFactory.java
  7. 2
      common/message/src/main/java/org/thingsboard/server/common/msg/stats/StatsType.java
  8. 77
      dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraBufferedRateExecutor.java
  9. 16
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java
  10. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueParams.java
  11. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java
  12. 6
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
  13. 5
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java
  14. 5
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
  15. 48
      dao/src/main/java/org/thingsboard/server/dao/util/AbstractBufferedRateExecutor.java
  16. 91
      dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateExecutorStats.java

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

22
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();
}
}

16
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 extends Number> 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);

9
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();
}

25
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;
}

4
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 extends Number> T createGauge(String key, T number, String... tags);
MessagesStats createMessagesStats(String key);
}

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

77
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

16
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<E> implements TbSqlQueue<E> {
private final BlockingQueue<TbSqlQueueElement<E>> 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<E> implements TbSqlQueue<E> {
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<E> implements TbSqlQueue<E> {
}
}
} 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<E> implements TbSqlQueue<E> {
});
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<E> implements TbSqlQueue<E> {
public ListenableFuture<Void> add(E element) {
SettableFuture<Void> future = SettableFuture.create();
queue.add(new TbSqlQueueElement<>(future, element));
addedCount.incrementAndGet();
params.getStats().incrementTotal();
return future;
}
}

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

5
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<AttributeKvEntity, Integer> hashcodeFunction = entity -> entity.getId().getEntityId().hashCode();

6
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<TsKvEntity> insertRepository;
protected TbSqlBlockingQueueWrapper<TsKvEntity> 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<TsKvEntity, Integer> hashcodeFunction = entity -> entity.getEntityId().hashCode();

5
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<TsKvLatestEntity, Integer> hashcodeFunction = entity -> entity.getEntityId().hashCode();

5
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<TimescaleTsKvEntity> insertRepository;
@ -74,6 +78,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
.batchSize(tsBatchSize)
.maxDelay(tsMaxDelay)
.statsPrintIntervalMs(tsStatsPrintIntervalMs)
.stats(statsFactory.createMessagesStats("ts.timescale"))
.build();
Function<TimescaleTsKvEntity, Integer> hashcodeFunction = entity -> entity.getEntityId().hashCode();

48
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<T extends AsyncTask, F extends ListenableFuture<V>, V> implements BufferedRateExecutor<T, F> {
public static final String CONCURRENCY_LEVEL = "currBuffer";
private final long maxWaitTime;
private final long pollMs;
private final BlockingQueue<AsyncTaskContext<T, V>> queue;
@ -49,20 +57,14 @@ public abstract class AbstractBufferedRateExecutor<T extends AsyncTask, F extend
private final boolean perTenantLimitsEnabled;
private final String perTenantLimitsConfiguration;
private final ConcurrentMap<TenantId, TbRateLimits> perTenantLimits = new ConcurrentHashMap<>();
protected final ConcurrentMap<TenantId, AtomicInteger> 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<T extends AsyncTask, F extend
this.timeoutExecutor = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("nosql-timeout"));
this.perTenantLimitsEnabled = perTenantLimitsEnabled;
this.perTenantLimitsConfiguration = perTenantLimitsConfiguration;
this.stats = new BufferedRateExecutorStats(statsFactory);
String concurrencyLevelKey = StatsType.RATE_EXECUTOR.getName() + "." + CONCURRENCY_LEVEL;
this.concurrencyLevel = statsFactory.createGauge(concurrencyLevelKey, new AtomicInteger(0));
for (int i = 0; i < dispatcherThreads; i++) {
dispatcherExecutor.submit(this::dispatch);
}
@ -89,8 +95,8 @@ public abstract class AbstractBufferedRateExecutor<T extends AsyncTask, F extend
} else if (!task.getTenantId().isNullUid()) {
TbRateLimits rateLimits = perTenantLimits.computeIfAbsent(task.getTenantId(), id -> 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<T extends AsyncTask, F extend
}
if (!perTenantLimitReached) {
try {
totalAdded.incrementAndGet();
stats.getTotalAdded().increment();
queue.add(new AsyncTaskContext<>(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<T extends AsyncTask, F extend
concurrencyLevel.incrementAndGet();
long timeout = finalTaskCtx.getCreateTime() + maxWaitTime - System.currentTimeMillis();
if (timeout > 0) {
totalLaunched.incrementAndGet();
stats.getTotalLaunched().increment();
ListenableFuture<V> result = execute(finalTaskCtx);
result = Futures.withTimeout(result, timeout, TimeUnit.MILLISECONDS, timeoutExecutor);
Futures.addCallback(result, new FutureCallback<V>() {
@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<T extends AsyncTask, F extend
} else {
logTask("Failed", finalTaskCtx);
}
totalFailed.incrementAndGet();
stats.getTotalFailed().increment();
concurrencyLevel.decrementAndGet();
finalTaskCtx.getFuture().setException(t);
log.debug("[{}] Failed to execute task: {}", finalTaskCtx.getId(), finalTaskCtx.getTask(), t);
@ -173,7 +179,7 @@ public abstract class AbstractBufferedRateExecutor<T extends AsyncTask, F extend
}, callbackExecutor);
} else {
logTask("Expired Before Execution", finalTaskCtx);
totalExpired.incrementAndGet();
stats.getTotalExpired().increment();
concurrencyLevel.decrementAndGet();
taskCtx.getFuture().setException(new TimeoutException());
}
@ -185,7 +191,7 @@ public abstract class AbstractBufferedRateExecutor<T extends AsyncTask, F extend
} catch (Throwable e) {
if (taskCtx != null) {
log.debug("[{}] Failed to execute task: {}", taskCtx.getId(), taskCtx, e);
totalFailed.incrementAndGet();
stats.getTotalFailed().increment();
concurrencyLevel.decrementAndGet();
} else {
log.debug("Failed to queue task:", e);

91
dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateExecutorStats.java

@ -0,0 +1,91 @@
/**
* 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.dao.util;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
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 java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicInteger;
@Slf4j
@Getter
public class BufferedRateExecutorStats {
private static final String TENANT_ID_TAG = "tenantId";
private static final String TOTAL_ADDED = "totalAdded";
private static final String TOTAL_LAUNCHED = "totalLaunched";
private static final String TOTAL_RELEASED = "totalReleased";
private static final String TOTAL_FAILED = "totalFailed";
private static final String TOTAL_EXPIRED = "totalExpired";
private static final String TOTAL_REJECTED = "totalRejected";
private static final String TOTAL_RATE_LIMITED = "totalRateLimited";
private final StatsFactory statsFactory;
private final ConcurrentMap<TenantId, DefaultCounter> rateLimitedTenants = new ConcurrentHashMap<>();
private final List<StatsCounter> 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();
}
}
Loading…
Cancel
Save