From 6809d374a6e61be66af2356c33722233a9aab403 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 5 Oct 2023 09:56:23 +0200 Subject: [PATCH 1/4] ACTORS_SYSTEM_RULE_DISPATCHER_POOL_SIZE increased from 4 to 8 --- .../thingsboard/server/actors/service/DefaultActorService.java | 2 +- application/src/main/resources/thingsboard.yml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java index 380da71443..87831cb8f6 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java @@ -74,7 +74,7 @@ public class DefaultActorService extends TbApplicationEventListener Date: Thu, 5 Oct 2023 10:41:29 +0200 Subject: [PATCH 2/4] SQL_ATTRIBUTES_BATCH_SIZE decreased from 10k down to 1k to improve concurrent updates on maximum load --- application/src/main/resources/thingsboard.yml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 0b012ba96c..9827a58a53 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -268,8 +268,8 @@ cassandra: sql: # Specify batch size for persisting attribute updates attributes: - batch_size: "${SQL_ATTRIBUTES_BATCH_SIZE:10000}" - batch_max_delay: "${SQL_ATTRIBUTES_BATCH_MAX_DELAY_MS:100}" + batch_size: "${SQL_ATTRIBUTES_BATCH_SIZE:1000}" + batch_max_delay: "${SQL_ATTRIBUTES_BATCH_MAX_DELAY_MS:50}" stats_print_interval_ms: "${SQL_ATTRIBUTES_BATCH_STATS_PRINT_MS:10000}" batch_threads: "${SQL_ATTRIBUTES_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution value_no_xss_validation: "${SQL_ATTRIBUTES_VALUE_NO_XSS_VALIDATION:false}" From c57f0bb8aa629cfcd260f1070db41d66c56bea12 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 5 Oct 2023 10:41:43 +0200 Subject: [PATCH 3/4] SQL_TS_LATEST_BATCH_SIZE decreased from 10k down to 1k to improve concurrent updates on maximum load --- application/src/main/resources/thingsboard.yml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 9827a58a53..4a3023d48f 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -280,8 +280,8 @@ sql: batch_threads: "${SQL_TS_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution value_no_xss_validation: "${SQL_TS_VALUE_NO_XSS_VALIDATION:false}" ts_latest: - batch_size: "${SQL_TS_LATEST_BATCH_SIZE:10000}" - batch_max_delay: "${SQL_TS_LATEST_BATCH_MAX_DELAY_MS:100}" + batch_size: "${SQL_TS_LATEST_BATCH_SIZE:1000}" + batch_max_delay: "${SQL_TS_LATEST_BATCH_MAX_DELAY_MS:50}" stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}" batch_threads: "${SQL_TS_LATEST_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution update_by_latest_ts: "${SQL_TS_UPDATE_BY_LATEST_TIMESTAMP:true}" From 31f3c2a824a534c830aefc6d8122884859f5a2b6 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 5 Oct 2023 11:39:11 +0200 Subject: [PATCH 4/4] CachedAttributesService find() many refactored in non-blocking manner --- .../dao/attributes/CachedAttributesService.java | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java index faff81670b..9a1a1b7ee3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java @@ -36,6 +36,7 @@ import org.thingsboard.server.common.stats.DefaultCounter; import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.dao.cache.CacheExecutorService; import org.thingsboard.server.dao.service.Validator; +import org.thingsboard.server.dao.sql.JpaExecutorService; import javax.annotation.PostConstruct; import java.util.ArrayList; @@ -61,6 +62,7 @@ public class CachedAttributesService implements AttributesService { public static final String LOCAL_CACHE_TYPE = "caffeine"; private final AttributesDao attributesDao; + private final JpaExecutorService jpaExecutorService; private final CacheExecutorService cacheExecutorService; private final DefaultCounter hitCounter; private final DefaultCounter missCounter; @@ -73,10 +75,12 @@ public class CachedAttributesService implements AttributesService { private boolean valueNoXssValidation; public CachedAttributesService(AttributesDao attributesDao, + JpaExecutorService jpaExecutorService, StatsFactory statsFactory, CacheExecutorService cacheExecutorService, TbTransactionalCache cache) { this.attributesDao = attributesDao; + this.jpaExecutorService = jpaExecutorService; this.cacheExecutorService = cacheExecutorService; this.cache = cache; @@ -134,12 +138,14 @@ public class CachedAttributesService implements AttributesService { } @Override - public ListenableFuture> find(TenantId tenantId, EntityId entityId, String scope, Collection attributeKeys) { + public ListenableFuture> find(TenantId tenantId, EntityId entityId, String scope, final Collection attributeKeysNonUnique) { validate(entityId, scope); - attributeKeys = new LinkedHashSet<>(attributeKeys); // deduplicate the attributes + final var attributeKeys = new LinkedHashSet<>(attributeKeysNonUnique); // deduplicate the attributes attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey)); - Map> wrappedCachedAttributes = findCachedAttributes(entityId, scope, attributeKeys); + //CacheExecutor for Redis or DirectExecutor for local Caffeine + return Futures.transformAsync(cacheExecutor.submit(() -> findCachedAttributes(entityId, scope, attributeKeys)), + wrappedCachedAttributes -> { List cachedAttributes = wrappedCachedAttributes.values().stream() .map(TbCacheValueWrapper::get) @@ -155,7 +161,8 @@ public class CachedAttributesService implements AttributesService { List notFoundKeys = notFoundAttributeKeys.stream().map(k -> new AttributeCacheKey(scope, entityId, k)).collect(Collectors.toList()); - return cacheExecutor.submit(() -> { + // DB call should run in DB executor, not in cache-related executor + return jpaExecutorService.submit(() -> { var cacheTransaction = cache.newTransactionForKeys(notFoundKeys); try { log.trace("[{}][{}] Lookup attributes from db: {}", entityId, scope, notFoundAttributeKeys); @@ -179,6 +186,8 @@ public class CachedAttributesService implements AttributesService { throw e; } }); + + }, MoreExecutors.directExecutor()); // cacheExecutor analyse and returns results or submit to DB executor } private Map> findCachedAttributes(EntityId entityId, String scope, Collection attributeKeys) {