From 71bd6001bab6d3be2b72c744bff0b97f8ee2a93d Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Thu, 25 Jun 2020 11:43:42 +0300 Subject: [PATCH 1/4] changed validateEmail access modifier to public --- .../java/org/thingsboard/server/dao/service/DataValidator.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java index 7f04ddf2aa..b6cca853bb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java @@ -64,7 +64,7 @@ public abstract class DataValidator> { return actualData.getId() != null && existentData.getId().equals(actualData.getId()); } - protected static void validateEmail(String email) { + public static void validateEmail(String email) { if (!doValidateEmail(email)) { throw new DataValidationException("Invalid email address format '" + email + "'!"); } From e77127019f058a00f7c5608b565c8d7ca12d2ae7 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 10 Jul 2020 12:46:47 +0300 Subject: [PATCH 2/4] Ts batch improvements --- .../src/main/resources/thingsboard.yml | 4 ++ .../server/dao/sql/TbSqlBlockingQueue.java | 8 ++- .../dao/sql/TbSqlBlockingQueueWrapper.java | 53 +++++++++++++++++++ .../server/dao/sql/TbSqlQueue.java | 2 +- .../dao/sql/attributes/JpaAttributeDao.java | 12 +++-- ...stractChunkedAggregationTimeseriesDao.java | 10 ++-- .../dao/sqlts/AbstractSqlTimeseriesDao.java | 34 ++++++++++-- .../thingsboard/server/dao/sqlts/TsKey.java | 26 +++++++++ .../timescale/TimescaleTimeseriesDao.java | 10 ++-- 9 files changed, 140 insertions(+), 19 deletions(-) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/TsKey.java diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index d25c94e49a..337df7ef12 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -252,14 +252,17 @@ sql: batch_size: "${SQL_ATTRIBUTES_BATCH_SIZE:10000}" batch_max_delay: "${SQL_ATTRIBUTES_BATCH_MAX_DELAY_MS:100}" stats_print_interval_ms: "${SQL_ATTRIBUTES_BATCH_STATS_PRINT_MS:10000}" + batch_threads: "${SQL_ATTRIBUTES_BATCH_THREADS:4}" ts: batch_size: "${SQL_TS_BATCH_SIZE:10000}" batch_max_delay: "${SQL_TS_BATCH_MAX_DELAY_MS:100}" stats_print_interval_ms: "${SQL_TS_BATCH_STATS_PRINT_MS:10000}" + batch_threads: "${SQL_TS_BATCH_THREADS:4}" ts_latest: batch_size: "${SQL_TS_LATEST_BATCH_SIZE:10000}" batch_max_delay: "${SQL_TS_LATEST_BATCH_MAX_DELAY_MS:100}" stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}" + batch_threads: "${SQL_TS_LATEST_BATCH_THREADS:4}" # Specify whether to remove null characters from strValue of attributes and timeseries before insert remove_null_chars: "${SQL_REMOVE_NULL_CHARS:true}" postgres: @@ -268,6 +271,7 @@ sql: timescale: # Specify Interval size for new data chunks storage. chunk_time_interval: "${SQL_TIMESCALE_CHUNK_TIME_INTERVAL:604800000}" + batch_threads: "${SQL_TIMESCALE_BATCH_THREADS:4}" ttl: ts: enabled: "${SQL_TTL_TS_ENABLED:true}" 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 e20a839e0e..3554fe7bce 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 @@ -41,16 +41,14 @@ public class TbSqlBlockingQueue implements TbSqlQueue { private final TbSqlBlockingQueueParams params; private ExecutorService executor; - private ScheduledLogExecutorComponent logExecutor; public TbSqlBlockingQueue(TbSqlBlockingQueueParams params) { this.params = params; } @Override - public void init(ScheduledLogExecutorComponent logExecutor, Consumer> saveFunction) { - this.logExecutor = logExecutor; - executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("sql-queue-" + params.getLogName().toLowerCase())); + public void init(ScheduledLogExecutorComponent logExecutor, Consumer> saveFunction, int index) { + executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("sql-queue-" + index + "-" + params.getLogName().toLowerCase())); executor.submit(() -> { String logName = params.getLogName(); int batchSize = params.getBatchSize(); @@ -94,7 +92,7 @@ public class TbSqlBlockingQueue implements TbSqlQueue { logExecutor.scheduleAtFixedRate(() -> { if (queue.size() > 0 || addedCount.get() > 0 || savedCount.get() > 0 || failedCount.get() > 0) { - log.info("[{}] queueSize [{}] totalAdded [{}] totalSaved [{}] totalFailed [{}]", + log.info("Queue-{} [{}] queueSize [{}] totalAdded [{}] totalSaved [{}] totalFailed [{}]", index, params.getLogName(), queue.size(), addedCount.getAndSet(0), savedCount.getAndSet(0), failedCount.getAndSet(0)); } }, params.getStatsPrintIntervalMs(), params.getStatsPrintIntervalMs(), TimeUnit.MILLISECONDS); 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 new file mode 100644 index 0000000000..1ad2b280e0 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java @@ -0,0 +1,53 @@ +/** + * 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.sql; + +import com.google.common.util.concurrent.ListenableFuture; +import lombok.Data; +import lombok.extern.slf4j.Slf4j; + +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.function.Consumer; +import java.util.function.Function; + +@Slf4j +@Data +public class TbSqlBlockingQueueWrapper { + private final CopyOnWriteArrayList> queues = new CopyOnWriteArrayList<>(); + private final TbSqlBlockingQueueParams params; + private ScheduledLogExecutorComponent logExecutor; + private final Function hashCodeFunction; + private final int maxThreads; + + public void init(ScheduledLogExecutorComponent logExecutor, Consumer> saveFunction) { + for (int i = 0; i < maxThreads; i++) { + TbSqlBlockingQueue queue = new TbSqlBlockingQueue<>(params); + queues.add(queue); + queue.init(logExecutor, saveFunction, i); + } + } + + public ListenableFuture add(E element) { + int hash = hashCodeFunction.apply(element); + int queueIndex = (hash & 0x7FFFFFFF) % maxThreads; + return queues.get(queueIndex).add(element); + } + + public void destroy() { + queues.forEach(TbSqlBlockingQueue::destroy); + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlQueue.java b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlQueue.java index d02c68a2d2..c3955a811c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlQueue.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlQueue.java @@ -22,7 +22,7 @@ import java.util.function.Consumer; public interface TbSqlQueue { - void init(ScheduledLogExecutorComponent logExecutor, Consumer> saveFunction); + void init(ScheduledLogExecutorComponent logExecutor, Consumer> saveFunction, int queueIndex); void destroy(); 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 c14e2dd7d0..56340d8729 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 @@ -32,8 +32,8 @@ import org.thingsboard.server.dao.model.sql.AttributeKvCompositeKey; import org.thingsboard.server.dao.model.sql.AttributeKvEntity; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; -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.util.SqlDao; import javax.annotation.PostConstruct; @@ -41,6 +41,7 @@ import javax.annotation.PreDestroy; import java.util.Collection; import java.util.List; import java.util.Optional; +import java.util.function.Function; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.UUIDConverter.fromTimeUUID; @@ -68,7 +69,10 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl @Value("${sql.attributes.stats_print_interval_ms:1000}") private long statsPrintIntervalMs; - private TbSqlBlockingQueue queue; + @Value("${sql.attributes.batch_threads:4}") + private int batchThreads; + + private TbSqlBlockingQueueWrapper queue; @PostConstruct private void init() { @@ -78,7 +82,9 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl .maxDelay(maxDelay) .statsPrintIntervalMs(statsPrintIntervalMs) .build(); - queue = new TbSqlBlockingQueue<>(params); + + Function hashcodeFunction = entity -> entity != null ? entity.getId().getEntityId().hashCode() : 0; + queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads); 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 c663d4346a..92b5b7e927 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 @@ -17,6 +17,7 @@ package org.thingsboard.server.dao.sqlts; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; @@ -31,8 +32,8 @@ import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; 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; import org.thingsboard.server.dao.sqlts.ts.TsKvRepository; import org.thingsboard.server.dao.timeseries.TimeseriesDao; @@ -43,6 +44,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 @@ -54,7 +56,7 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq @Autowired protected InsertTsRepository insertRepository; - protected TbSqlBlockingQueue tsQueue; + protected TbSqlBlockingQueueWrapper tsQueue; @PostConstruct protected void init() { @@ -65,7 +67,9 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq .maxDelay(tsMaxDelay) .statsPrintIntervalMs(tsStatsPrintIntervalMs) .build(); - tsQueue = new TbSqlBlockingQueue<>(tsParams); + + Function hashcodeFunction = entity -> entity != null ? entity.getEntityId().hashCode() : 0; + tsQueue = new TbSqlBlockingQueueWrapper<>(tsParams, hashcodeFunction, tsBatchThreads); 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 1d97aaddd8..81ef7e100a 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 @@ -41,8 +41,8 @@ import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestCompositeKey; import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; -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.dictionary.TsKvDictionaryRepository; import org.thingsboard.server.dao.sqlts.insert.latest.InsertLatestTsRepository; import org.thingsboard.server.dao.sqlts.latest.SearchTsKvLatestRepository; @@ -52,7 +52,9 @@ import org.thingsboard.server.dao.timeseries.SimpleListenableFuture; import javax.annotation.Nullable; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; +import java.util.Comparator; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; @@ -82,7 +84,7 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx @Autowired private TsKvDictionaryRepository dictionaryRepository; - private TbSqlBlockingQueue tsLatestQueue; + private TbSqlBlockingQueueWrapper tsLatestQueue; @Value("${sql.ts_latest.batch_size:1000}") private int tsLatestBatchSize; @@ -93,6 +95,9 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx @Value("${sql.ts_latest.stats_print_interval_ms:1000}") private long tsLatestStatsPrintIntervalMs; + @Value("${sql.ts_latest.batch_threads:4}") + private int tsLatestBatchThreads; + @Autowired protected ScheduledLogExecutorComponent logExecutor; @@ -105,6 +110,12 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx @Value("${sql.ts.stats_print_interval_ms:1000}") protected long tsStatsPrintIntervalMs; + @Value("${sql.ts.batch_threads:4}") + protected int tsBatchThreads; + + @Value("${sql.timescale.batch_threads:4}") + protected int timescaleBatchThreads; + @PostConstruct protected void init() { TbSqlBlockingQueueParams tsLatestParams = TbSqlBlockingQueueParams.builder() @@ -113,8 +124,23 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx .maxDelay(tsLatestMaxDelay) .statsPrintIntervalMs(tsLatestStatsPrintIntervalMs) .build(); - tsLatestQueue = new TbSqlBlockingQueue<>(tsLatestParams); - tsLatestQueue.init(logExecutor, v -> insertLatestTsRepository.saveOrUpdate(v)); + + java.util.function.Function hashcodeFunction = entity -> entity != null ? entity.getEntityId().hashCode() : 0; + tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads); + + tsLatestQueue.init(logExecutor, v -> { + Map> tsMap = + v.stream().collect(Collectors.groupingBy(ts -> new TsKey(ts.getEntityId(), ts.getStrKey()))); + + List latestEntities = + tsMap.keySet() + .stream() + .map(tsMap::get) + .map(list -> list.stream().max(Comparator.comparing(TsKvLatestEntity::getTs)).get()) + .collect(Collectors.toList()); + + insertLatestTsRepository.saveOrUpdate(latestEntities); + }); } @PreDestroy diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/TsKey.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/TsKey.java new file mode 100644 index 0000000000..d14615ca48 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/TsKey.java @@ -0,0 +1,26 @@ +/** + * 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.sqlts; + +import lombok.Data; + +import java.util.UUID; + +@Data +public class TsKey { + private final UUID entityId; + private final String key; +} 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 d9121f684e..04bbe58a14 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 @@ -33,8 +33,8 @@ import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.model.sqlts.timescale.ts.TimescaleTsKvEntity; -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.AbstractSqlTimeseriesDao; import org.thingsboard.server.dao.sqlts.insert.InsertTsRepository; import org.thingsboard.server.dao.timeseries.TimeseriesDao; @@ -48,6 +48,7 @@ import java.util.List; import java.util.Optional; import java.util.UUID; import java.util.concurrent.CompletableFuture; +import java.util.function.Function; @Component @Slf4j @@ -63,7 +64,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements @Autowired protected InsertTsRepository insertRepository; - protected TbSqlBlockingQueue tsQueue; + protected TbSqlBlockingQueueWrapper tsQueue; @PostConstruct protected void init() { @@ -74,7 +75,10 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements .maxDelay(tsMaxDelay) .statsPrintIntervalMs(tsStatsPrintIntervalMs) .build(); - tsQueue = new TbSqlBlockingQueue<>(tsParams); + + Function hashcodeFunction = entity -> entity != null ? entity.getEntityId().hashCode() : 0; + tsQueue = new TbSqlBlockingQueueWrapper<>(tsParams, hashcodeFunction, timescaleBatchThreads); + tsQueue.init(logExecutor, v -> insertRepository.saveOrUpdate(v)); } From 98c24632e3e87ac261adf581105d4bab75122ece Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 10 Jul 2020 15:24:11 +0300 Subject: [PATCH 3/4] refactored --- .../java/org/thingsboard/server/dao/service/DataValidator.java | 2 +- .../thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java | 3 +-- .../thingsboard/server/dao/sql/attributes/JpaAttributeDao.java | 2 +- .../dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java | 2 +- .../thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java | 2 +- .../server/dao/sqlts/timescale/TimescaleTimeseriesDao.java | 3 ++- 6 files changed, 7 insertions(+), 7 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java index b6cca853bb..7f04ddf2aa 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java @@ -64,7 +64,7 @@ public abstract class DataValidator> { return actualData.getId() != null && existentData.getId().equals(actualData.getId()); } - public static void validateEmail(String email) { + protected static void validateEmail(String email) { if (!doValidateEmail(email)) { throw new DataValidationException("Invalid email address format '" + email + "'!"); } 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 1ad2b280e0..2c53cc4ad9 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 @@ -42,8 +42,7 @@ public class TbSqlBlockingQueueWrapper { } public ListenableFuture add(E element) { - int hash = hashCodeFunction.apply(element); - int queueIndex = (hash & 0x7FFFFFFF) % maxThreads; + int queueIndex = element != null ? (hashCodeFunction.apply(element) & 0x7FFFFFFF) % maxThreads : 0; return queues.get(queueIndex).add(element); } 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 56340d8729..a3b17c7863 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 @@ -83,7 +83,7 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl .statsPrintIntervalMs(statsPrintIntervalMs) .build(); - Function hashcodeFunction = entity -> entity != null ? entity.getId().getEntityId().hashCode() : 0; + Function hashcodeFunction = entity -> entity.getId().getEntityId().hashCode(); queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads); 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 92b5b7e927..84dddcaced 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 @@ -68,7 +68,7 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq .statsPrintIntervalMs(tsStatsPrintIntervalMs) .build(); - Function hashcodeFunction = entity -> entity != null ? entity.getEntityId().hashCode() : 0; + Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); tsQueue = new TbSqlBlockingQueueWrapper<>(tsParams, hashcodeFunction, tsBatchThreads); 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 81ef7e100a..0e2dfabbef 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 @@ -125,7 +125,7 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx .statsPrintIntervalMs(tsLatestStatsPrintIntervalMs) .build(); - java.util.function.Function hashcodeFunction = entity -> entity != null ? entity.getEntityId().hashCode() : 0; + java.util.function.Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads); tsLatestQueue.init(logExecutor, v -> { 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 04bbe58a14..a84a3e884f 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 @@ -76,7 +76,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements .statsPrintIntervalMs(tsStatsPrintIntervalMs) .build(); - Function hashcodeFunction = entity -> entity != null ? entity.getEntityId().hashCode() : 0; + Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); tsQueue = new TbSqlBlockingQueueWrapper<>(tsParams, hashcodeFunction, timescaleBatchThreads); tsQueue.init(logExecutor, v -> insertRepository.saveOrUpdate(v)); @@ -281,4 +281,5 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements startTs, endTs); } + } From 2f4ca3b5be4cbeea708c22b0553065d859b32743 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 10 Jul 2020 17:07:15 +0300 Subject: [PATCH 4/4] improvements --- .../dao/sqlts/AbstractSqlTimeseriesDao.java | 22 +++++++++---------- .../thingsboard/server/dao/sqlts/TsKey.java | 2 +- 2 files changed, 12 insertions(+), 12 deletions(-) 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 0e2dfabbef..75b3d6951b 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 @@ -52,7 +52,8 @@ import org.thingsboard.server.dao.timeseries.SimpleListenableFuture; import javax.annotation.Nullable; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; -import java.util.Comparator; +import java.util.ArrayList; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Objects; @@ -129,16 +130,15 @@ public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningEx tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads); tsLatestQueue.init(logExecutor, v -> { - Map> tsMap = - v.stream().collect(Collectors.groupingBy(ts -> new TsKey(ts.getEntityId(), ts.getStrKey()))); - - List latestEntities = - tsMap.keySet() - .stream() - .map(tsMap::get) - .map(list -> list.stream().max(Comparator.comparing(TsKvLatestEntity::getTs)).get()) - .collect(Collectors.toList()); - + Map trueLatest = new HashMap<>(); + v.forEach(ts -> { + TsKey key = new TsKey(ts.getEntityId(), ts.getKey()); + TsKvLatestEntity old = trueLatest.get(key); + if (old == null || old.getTs() < ts.getTs()) { + trueLatest.put(key, ts); + } + }); + List latestEntities = new ArrayList<>(trueLatest.values()); insertLatestTsRepository.saveOrUpdate(latestEntities); }); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/TsKey.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/TsKey.java index d14615ca48..17da2f80bc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/TsKey.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/TsKey.java @@ -22,5 +22,5 @@ import java.util.UUID; @Data public class TsKey { private final UUID entityId; - private final String key; + private final int key; }