Browse Source

Merged with develop/2.5.3

pull/3199/head
vzikratyi 6 years ago
committed by Andrew Shvayka
parent
commit
5b4cfbb3d8
  1. 12
      common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java
  2. 7
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java
  3. 3
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueParams.java
  4. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java
  5. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java
  6. 5
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
  7. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java
  8. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java

12
common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java

@ -174,17 +174,23 @@ public class DefaultTbQueueRequestTemplate<Request extends TbQueueMsg, Response
ResponseMetaData<Response> 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);
}

7
dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java

@ -38,10 +38,11 @@ public class TbSqlBlockingQueue<E> implements TbSqlQueue<E> {
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<E> implements TbSqlQueue<E> {
public ListenableFuture<Void> add(E element) {
SettableFuture<Void> future = SettableFuture.create();
queue.add(new TbSqlQueueElement<>(future, element));
params.getStats().incrementTotal();
stats.incrementTotal();
return future;
}
}

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

6
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<E> {
private ScheduledLogExecutorComponent logExecutor;
private final Function<E, Integer> hashCodeFunction;
private final int maxThreads;
private final StatsFactory statsFactory;
public void init(ScheduledLogExecutorComponent logExecutor, Consumer<List<E>> saveFunction) {
for (int i = 0; i < maxThreads; i++) {
TbSqlBlockingQueue<E> queue = new TbSqlBlockingQueue<>(params);
MessagesStats stats = statsFactory.createMessagesStats(params.getStatsNamePrefix() + ".queue." + i);
TbSqlBlockingQueue<E> queue = new TbSqlBlockingQueue<>(params, stats);
queues.add(queue);
queue.init(logExecutor, saveFunction, i);
}

4
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<AttributeKvEntity, Integer> 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));
}

5
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<TsKvEntity, Integer> 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));
}

4
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<TsKvLatestEntity, Integer> hashcodeFunction = entity -> entity.getEntityId().hashCode();
tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads);
tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads, statsFactory);
tsLatestQueue.init(logExecutor, v -> {
Map<TsKey, TsKvLatestEntity> trueLatest = new HashMap<>();

4
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<TimescaleTsKvEntity, Integer> 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));
}

Loading…
Cancel
Save