From 3ac0cc241a23af1ecc885055f81f2d4bbcff0f2d Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Wed, 27 Dec 2023 13:47:52 +0200 Subject: [PATCH 01/11] TbAwsSqsProducerTemplate: native thread allocation error fix (provided one executor for message sending with limited thread pool size) --- .../queue/sqs/TbAwsSqsProducerTemplate.java | 38 +++++++------------ .../server/queue/sqs/TbAwsSqsSettings.java | 18 +++++++++ 2 files changed, 31 insertions(+), 25 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java index 863e96a714..ab602e1960 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java @@ -20,15 +20,11 @@ import com.amazonaws.auth.AWSCredentialsProvider; import com.amazonaws.auth.AWSStaticCredentialsProvider; import com.amazonaws.auth.BasicAWSCredentials; import com.amazonaws.auth.DefaultAWSCredentialsProviderChain; -import com.amazonaws.services.sqs.AmazonSQS; -import com.amazonaws.services.sqs.AmazonSQSClientBuilder; +import com.amazonaws.handlers.AsyncHandler; +import com.amazonaws.services.sqs.AmazonSQSAsync; +import com.amazonaws.services.sqs.AmazonSQSAsyncClientBuilder; import com.amazonaws.services.sqs.model.SendMessageRequest; import com.amazonaws.services.sqs.model.SendMessageResult; -import com.google.common.util.concurrent.FutureCallback; -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.gson.Gson; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -41,16 +37,14 @@ import org.thingsboard.server.queue.common.DefaultTbQueueMsg; import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; @Slf4j public class TbAwsSqsProducerTemplate implements TbQueueProducer { private final String defaultTopic; - private final AmazonSQS sqsClient; + private final AmazonSQSAsync sqsClient; private final Gson gson = new Gson(); private final Map queueUrlMap = new ConcurrentHashMap<>(); private final TbQueueAdmin admin; - private ListeningExecutorService producerExecutor; public TbAwsSqsProducerTemplate(TbQueueAdmin admin, TbAwsSqsSettings sqsSettings, String defaultTopic) { this.admin = admin; @@ -64,11 +58,11 @@ public class TbAwsSqsProducerTemplate implements TbQueuePr credentialsProvider = new AWSStaticCredentialsProvider(awsCredentials); } - sqsClient = AmazonSQSClientBuilder.standard() + sqsClient = AmazonSQSAsyncClientBuilder.standard() .withCredentials(credentialsProvider) .withRegion(sqsSettings.getRegion()) + .withExecutorFactory(sqsSettings::getProducerExecutor) .build(); - producerExecutor = MoreExecutors.listeningDecorator(Executors.newCachedThreadPool()); } @Override @@ -91,30 +85,24 @@ public class TbAwsSqsProducerTemplate implements TbQueuePr sendMsgRequest.withMessageGroupId(sqsMsgId); sendMsgRequest.withMessageDeduplicationId(sqsMsgId); - ListenableFuture future = producerExecutor.submit(() -> sqsClient.sendMessage(sendMsgRequest)); - - Futures.addCallback(future, new FutureCallback() { - @Override - public void onSuccess(SendMessageResult result) { + sqsClient.sendMessageAsync(sendMsgRequest, new AsyncHandler() { + @Override public void onError(Exception e) { if (callback != null) { - callback.onSuccess(new AwsSqsTbQueueMsgMetadata(result.getSdkHttpMetadata())); + callback.onFailure(e); } } - @Override - public void onFailure(Throwable t) { + @Override public void onSuccess(SendMessageRequest request, + SendMessageResult sendMessageResult) { if (callback != null) { - callback.onFailure(t); + callback.onSuccess(new AwsSqsTbQueueMsgMetadata(sendMessageResult.getSdkHttpMetadata())); } } - }, producerExecutor); + }); } @Override public void stop() { - if (producerExecutor != null) { - producerExecutor.shutdownNow(); - } if (sqsClient != null) { sqsClient.shutdown(); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java index 89dc9d7e47..0a3915496b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java @@ -20,6 +20,11 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Component; +import org.thingsboard.common.util.ThingsBoardThreadFactory; + +import javax.annotation.PostConstruct; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; @Slf4j @ConditionalOnExpression("'${queue.type:null}'=='aws-sqs'") @@ -42,4 +47,17 @@ public class TbAwsSqsSettings { @Value("${queue.aws_sqs.threads_per_topic}") private int threadsPerTopic; + @Value("${queue.aws_sqs.producer_thread_pool_size:0}") + private int threadPoolSize; + + private ExecutorService producerExecutor; + + @PostConstruct + private void init() { + if (threadPoolSize == 0) { + threadPoolSize = 50; //AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE = 50; + } + producerExecutor = Executors.newFixedThreadPool(threadPoolSize, ThingsBoardThreadFactory.forName("aws-sqs-queue-executor")); + } + } From 63a85161b2980a71221dae2fcd68986ef3767976 Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Wed, 27 Dec 2023 14:44:49 +0200 Subject: [PATCH 02/11] moved producer executor to TbAwsSqsAdmin --- .../thingsboard/server/queue/sqs/TbAwsSqsAdmin.java | 11 +++++++++++ .../server/queue/sqs/TbAwsSqsProducerTemplate.java | 6 +++--- .../server/queue/sqs/TbAwsSqsSettings.java | 10 ---------- 3 files changed, 14 insertions(+), 13 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java index ba4eeb6ca4..337bb271e5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java @@ -24,11 +24,15 @@ import com.amazonaws.services.sqs.AmazonSQS; import com.amazonaws.services.sqs.AmazonSQSClientBuilder; import com.amazonaws.services.sqs.model.CreateQueueRequest; import com.amazonaws.services.sqs.model.GetQueueUrlResult; +import lombok.Getter; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.util.PropertyUtils; import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.function.Function; import java.util.stream.Collectors; @@ -38,6 +42,8 @@ public class TbAwsSqsAdmin implements TbQueueAdmin { private final Map attributes; private final AmazonSQS sqsClient; private final Map queues; + @Getter + private ExecutorService producerExecutor; public TbAwsSqsAdmin(TbAwsSqsSettings sqsSettings, Map attributes) { this.attributes = attributes; @@ -49,6 +55,11 @@ public class TbAwsSqsAdmin implements TbQueueAdmin { AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); credentialsProvider = new AWSStaticCredentialsProvider(awsCredentials); } + int threadPoolSize = sqsSettings.getThreadPoolSize(); + if (threadPoolSize == 0) { + threadPoolSize = 50; //AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE = 50; + } + producerExecutor = Executors.newFixedThreadPool(threadPoolSize, ThingsBoardThreadFactory.forName("aws-sqs-queue-executor")); sqsClient = AmazonSQSClientBuilder.standard() .withCredentials(credentialsProvider) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java index ab602e1960..f5bca724d0 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java @@ -44,10 +44,10 @@ public class TbAwsSqsProducerTemplate implements TbQueuePr private final AmazonSQSAsync sqsClient; private final Gson gson = new Gson(); private final Map queueUrlMap = new ConcurrentHashMap<>(); - private final TbQueueAdmin admin; + private final TbAwsSqsAdmin admin; public TbAwsSqsProducerTemplate(TbQueueAdmin admin, TbAwsSqsSettings sqsSettings, String defaultTopic) { - this.admin = admin; + this.admin = (TbAwsSqsAdmin) admin; this.defaultTopic = defaultTopic; AWSCredentialsProvider credentialsProvider; @@ -61,7 +61,7 @@ public class TbAwsSqsProducerTemplate implements TbQueuePr sqsClient = AmazonSQSAsyncClientBuilder.standard() .withCredentials(credentialsProvider) .withRegion(sqsSettings.getRegion()) - .withExecutorFactory(sqsSettings::getProducerExecutor) + .withExecutorFactory(this.admin::getProducerExecutor) .build(); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java index 0a3915496b..1f8e301c87 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java @@ -50,14 +50,4 @@ public class TbAwsSqsSettings { @Value("${queue.aws_sqs.producer_thread_pool_size:0}") private int threadPoolSize; - private ExecutorService producerExecutor; - - @PostConstruct - private void init() { - if (threadPoolSize == 0) { - threadPoolSize = 50; //AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE = 50; - } - producerExecutor = Executors.newFixedThreadPool(threadPoolSize, ThingsBoardThreadFactory.forName("aws-sqs-queue-executor")); - } - } From f2ea1b87d76328b9bd2b46656f504fe8db75de27 Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Wed, 27 Dec 2023 14:51:20 +0200 Subject: [PATCH 03/11] added shutdown for producerExecutor --- .../java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java index 337bb271e5..7ef3b26b27 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java @@ -115,5 +115,8 @@ public class TbAwsSqsAdmin implements TbQueueAdmin { if (sqsClient != null) { sqsClient.shutdown(); } + if (producerExecutor != null) { + producerExecutor.shutdownNow(); + } } } From b95adb439a69ca231d57ebfc4fe94afe6aebb9fe Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Wed, 27 Dec 2023 17:06:26 +0200 Subject: [PATCH 04/11] added yml parameters for aws-sqs eproducer executor thread pool size --- application/src/main/resources/thingsboard.yml | 2 ++ .../org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java | 8 ++------ .../thingsboard/server/queue/sqs/TbAwsSqsSettings.java | 7 +------ msa/vc-executor/src/main/resources/tb-vc-executor.yml | 2 ++ transport/coap/src/main/resources/tb-coap-transport.yml | 2 ++ transport/http/src/main/resources/tb-http-transport.yml | 2 ++ transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml | 2 ++ transport/mqtt/src/main/resources/tb-mqtt-transport.yml | 2 ++ transport/snmp/src/main/resources/tb-snmp-transport.yml | 2 ++ 9 files changed, 17 insertions(+), 12 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index c650d53b8f..0561c0c249 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1398,6 +1398,8 @@ queue: region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}" # Number of threads per each AWS SQS queue in consumer threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}" + # Thread pool size for aws_sqs queue producer executor provider. Default value equals to AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE + producer_thread_pool_size: "${TB_QUEUE_AWS_SQS_EXECUTOR_THREAD_POOL_SIZE:50}" queue-properties: # AWS SQS queue properties. VisibilityTimeout in seconds;MaximumMessageSize in bytes;MessageRetentionPeriod in seconds rule-engine: "${TB_QUEUE_AWS_SQS_RE_QUEUE_PROPERTIES:VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800}" diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java index 7ef3b26b27..1c70302d00 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java @@ -43,7 +43,7 @@ public class TbAwsSqsAdmin implements TbQueueAdmin { private final AmazonSQS sqsClient; private final Map queues; @Getter - private ExecutorService producerExecutor; + private final ExecutorService producerExecutor; public TbAwsSqsAdmin(TbAwsSqsSettings sqsSettings, Map attributes) { this.attributes = attributes; @@ -55,11 +55,7 @@ public class TbAwsSqsAdmin implements TbQueueAdmin { AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); credentialsProvider = new AWSStaticCredentialsProvider(awsCredentials); } - int threadPoolSize = sqsSettings.getThreadPoolSize(); - if (threadPoolSize == 0) { - threadPoolSize = 50; //AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE = 50; - } - producerExecutor = Executors.newFixedThreadPool(threadPoolSize, ThingsBoardThreadFactory.forName("aws-sqs-queue-executor")); + producerExecutor = Executors.newFixedThreadPool(sqsSettings.getThreadPoolSize(), ThingsBoardThreadFactory.forName("aws-sqs-queue-executor")); sqsClient = AmazonSQSClientBuilder.standard() .withCredentials(credentialsProvider) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java index 1f8e301c87..f2dc1efdfb 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java @@ -20,11 +20,6 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Component; -import org.thingsboard.common.util.ThingsBoardThreadFactory; - -import javax.annotation.PostConstruct; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; @Slf4j @ConditionalOnExpression("'${queue.type:null}'=='aws-sqs'") @@ -47,7 +42,7 @@ public class TbAwsSqsSettings { @Value("${queue.aws_sqs.threads_per_topic}") private int threadsPerTopic; - @Value("${queue.aws_sqs.producer_thread_pool_size:0}") + @Value("${queue.aws_sqs.producer_thread_pool_size:50}") private int threadPoolSize; } diff --git a/msa/vc-executor/src/main/resources/tb-vc-executor.yml b/msa/vc-executor/src/main/resources/tb-vc-executor.yml index 962c5303c7..e4519529c5 100644 --- a/msa/vc-executor/src/main/resources/tb-vc-executor.yml +++ b/msa/vc-executor/src/main/resources/tb-vc-executor.yml @@ -155,6 +155,8 @@ queue: region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}" # Number of threads per each AWS SQS queue in consumer threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}" + # Thread pool size for aws_sqs queue producer executor provider. Default value equals to AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE + producer_thread_pool_size: "${TB_QUEUE_AWS_SQS_EXECUTOR_THREAD_POOL_SIZE:50}" queue-properties: # AWS SQS queue properties. VisibilityTimeout in seconds;MaximumMessageSize in bytes;MessageRetentionPeriod in seconds core: "${TB_QUEUE_AWS_SQS_CORE_QUEUE_PROPERTIES:VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800}" diff --git a/transport/coap/src/main/resources/tb-coap-transport.yml b/transport/coap/src/main/resources/tb-coap-transport.yml index ab5deff3f7..f11b1b8ce2 100644 --- a/transport/coap/src/main/resources/tb-coap-transport.yml +++ b/transport/coap/src/main/resources/tb-coap-transport.yml @@ -281,6 +281,8 @@ queue: region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}" # Number of threads per each AWS SQS queue in consumer threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}" + # Thread pool size for aws_sqs queue producer executor provider. Default value equals to AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE + producer_thread_pool_size: "${TB_QUEUE_AWS_SQS_EXECUTOR_THREAD_POOL_SIZE:50}" queue-properties: # AWS SQS queue properties. VisibilityTimeout in seconds;MaximumMessageSize in bytes;MessageRetentionPeriod in seconds rule-engine: "${TB_QUEUE_AWS_SQS_RE_QUEUE_PROPERTIES:VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800}" diff --git a/transport/http/src/main/resources/tb-http-transport.yml b/transport/http/src/main/resources/tb-http-transport.yml index 74be81e667..b2fa93c434 100644 --- a/transport/http/src/main/resources/tb-http-transport.yml +++ b/transport/http/src/main/resources/tb-http-transport.yml @@ -265,6 +265,8 @@ queue: region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}" # Number of threads per each AWS SQS queue in consumer threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}" + # Thread pool size for aws_sqs queue producer executor provider. Default value equals to AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE + producer_thread_pool_size: "${TB_QUEUE_AWS_SQS_EXECUTOR_THREAD_POOL_SIZE:50}" queue-properties: # AWS SQS queue properties. VisibilityTimeout in seconds;MaximumMessageSize in bytes;MessageRetentionPeriod in seconds rule-engine: "${TB_QUEUE_AWS_SQS_RE_QUEUE_PROPERTIES:VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800}" diff --git a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml index 10dbac101c..c2891994b1 100644 --- a/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml +++ b/transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml @@ -360,6 +360,8 @@ queue: region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}" # Number of threads per each AWS SQS queue in consumer threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}" + # Thread pool size for aws_sqs queue producer executor provider. Default value equals to AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE + producer_thread_pool_size: "${TB_QUEUE_AWS_SQS_EXECUTOR_THREAD_POOL_SIZE:50}" queue-properties: # AWS SQS queue properties. VisibilityTimeout in seconds;MaximumMessageSize in bytes;MessageRetentionPeriod in seconds rule-engine: "${TB_QUEUE_AWS_SQS_RE_QUEUE_PROPERTIES:VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800}" diff --git a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml index 0eb8b4315b..519e7a8b32 100644 --- a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml +++ b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml @@ -297,6 +297,8 @@ queue: region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}" # Number of threads per each AWS SQS queue in consumer threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}" + # Thread pool size for aws_sqs queue producer executor provider. Default value equals to AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE + producer_thread_pool_size: "${TB_QUEUE_AWS_SQS_EXECUTOR_THREAD_POOL_SIZE:50}" queue-properties: # AWS SQS queue properties. VisibilityTimeout in seconds;MaximumMessageSize in bytes;MessageRetentionPeriod in seconds rule-engine: "${TB_QUEUE_AWS_SQS_RE_QUEUE_PROPERTIES:VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800}" diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index d139b641a6..ab9d6ad584 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/transport/snmp/src/main/resources/tb-snmp-transport.yml @@ -250,6 +250,8 @@ queue: region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}" # Number of threads per each AWS SQS queue in consumer threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}" + # Thread pool size for aws_sqs queue producer executor provider. Default value equals to AmazonSQSAsyncClient.DEFAULT_THREAD_POOL_SIZE + producer_thread_pool_size: "${TB_QUEUE_AWS_SQS_EXECUTOR_THREAD_POOL_SIZE:50}" queue-properties: # AWS SQS queue properties. VisibilityTimeout in seconds;MaximumMessageSize in bytes;MessageRetentionPeriod in seconds rule-engine: "${TB_QUEUE_AWS_SQS_RE_QUEUE_PROPERTIES:VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800}" From abfcf32c54e2705171d93abbb563fbdec97def5d Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Thu, 28 Dec 2023 17:55:56 +0200 Subject: [PATCH 05/11] replaced slow query with multi-call of indexed query --- .../update/DefaultDataUpdateService.java | 10 +++++++++- .../server/dao/rule/RuleChainService.java | 4 ++++ .../server/dao/rule/BaseRuleChainService.java | 10 ++++++++++ .../server/dao/rule/RuleNodeDao.java | 2 ++ .../server/dao/service/Validator.java | 15 ++++++++++++++- .../server/dao/sql/rule/JpaRuleNodeDao.java | 10 ++++++++++ .../dao/sql/rule/RuleNodeRepository.java | 7 ++++++- .../dao/sql/rule/JpaRuleNodeDaoTest.java | 19 +++++++++++++++++++ 8 files changed, 74 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java index 7d08fdca00..d0f45ac06c 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java @@ -313,9 +313,17 @@ public class DefaultDataUpdateService implements DataUpdateService { } private List getRuleNodesIdsWithTypeAndVersionLessThan(String type, int toVersion) { + var ruleNodeIds = new ArrayList(); + for (int v = toVersion - 1; v >= 0; v--) { + ruleNodeIds.addAll(getRuleNodesIdsWithTypeAndVersion(type, v)); + } + return ruleNodeIds; + } + + private List getRuleNodesIdsWithTypeAndVersion(String type, int version) { var ruleNodeIds = new ArrayList(); new PageDataIterable<>(pageLink -> - ruleChainService.findAllRuleNodeIdsByTypeAndVersionLessThan(type, toVersion, pageLink), DEFAULT_PAGE_SIZE + ruleChainService.findAllRuleNodeIdsByTypeAndVersion(type, version, pageLink), DEFAULT_PAGE_SIZE ).forEach(ruleNodeIds::add); return ruleNodeIds; } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java index 532da4ac00..03886bfecc 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java @@ -100,10 +100,14 @@ public interface RuleChainService extends EntityDaoService { PageData findAllRuleNodesByType(String type, PageLink pageLink); + @Deprecated(forRemoval = true, since = "3.6.3") PageData findAllRuleNodesByTypeAndVersionLessThan(String type, int version, PageLink pageLink); + @Deprecated(forRemoval = true, since = "3.6.3") PageData findAllRuleNodeIdsByTypeAndVersionLessThan(String type, int version, PageLink pageLink); + PageData findAllRuleNodeIdsByTypeAndVersion(String type, int version, PageLink pageLink); + List findAllRuleNodesByIds(List ruleNodeIds); RuleNode saveRuleNode(TenantId tenantId, RuleNode ruleNode); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 6257bb0420..5a6d0c1d74 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -81,6 +81,7 @@ import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateIds; import static org.thingsboard.server.dao.service.Validator.validatePageLink; import static org.thingsboard.server.dao.service.Validator.validatePositiveNumber; +import static org.thingsboard.server.dao.service.Validator.validateNonNegativeNumber; import static org.thingsboard.server.dao.service.Validator.validateString; /** @@ -737,6 +738,15 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC return ruleNodeDao.findAllRuleNodeIdsByTypeAndVersionLessThan(type, version, pageLink); } + @Override + public PageData findAllRuleNodeIdsByTypeAndVersion(String type, int version, PageLink pageLink) { + log.trace("Executing findAllRuleNodeIdsByTypeAndVersion, type {}, pageLink {}, version {}", type, pageLink, version); + validateString(type, "Incorrect type of the rule node"); + validateNonNegativeNumber(version, "Incorrect version. Version should be non-negative!"); + validatePageLink(pageLink); + return ruleNodeDao.findAllRuleNodeIdsByTypeAndVersion(type, version, pageLink); + } + @Override public List findAllRuleNodesByIds(List ruleNodeIds) { log.trace("Executing findAllRuleNodesByIds, ruleNodeIds {}", ruleNodeIds); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeDao.java b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeDao.java index b6f2dd097e..4927f2d68d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeDao.java @@ -38,6 +38,8 @@ public interface RuleNodeDao extends Dao { PageData findAllRuleNodeIdsByTypeAndVersionLessThan(String type, int version, PageLink pageLink); + PageData findAllRuleNodeIdsByTypeAndVersion(String type, int version, PageLink pageLink); + List findAllRuleNodeByIds(List ruleNodeIds); List findByExternalIds(RuleChainId ruleChainId, List externalIds); diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java b/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java index d8c4ffddc9..a712f7c775 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java @@ -61,7 +61,7 @@ public class Validator { /** - * This method validate long value. If value isn't possitive than throw + * This method validate long value. If value isn't positive than throw * IncorrectParameterException exception * * @param val the val @@ -73,6 +73,19 @@ public class Validator { } } + /** + * This method validate long value. If value is negative than throw + * IncorrectParameterException exception + * + * @param val the val + * @param errorMessage the error message for exception + */ + public static void validateNonNegativeNumber(long val, String errorMessage) { + if (val < 0) { + throw new IncorrectParameterException(errorMessage); + } + } + /** * This method validate UUID id. If id is null than throw * IncorrectParameterException exception diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java index de5ff109d3..655a03caf3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java @@ -89,6 +89,16 @@ public class JpaRuleNodeDao extends JpaAbstractDao imp .mapData(RuleNodeId::new); } + @Override + public PageData findAllRuleNodeIdsByTypeAndVersion(String type, int version, PageLink pageLink) { + return DaoUtil.pageToPageData(ruleNodeRepository + .findAllRuleNodeIdsByTypeAndVersion( + type, + version, + DaoUtil.toPageable(pageLink))) + .mapData(RuleNodeId::new); + } + @Override public List findAllRuleNodeByIds(List ruleNodeIds) { return DaoUtil.convertDataList(ruleNodeRepository.findAllById( diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java index f02072b356..eaccf39072 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java @@ -43,7 +43,7 @@ public interface RuleNodeRepository extends JpaRepository Pageable pageable); @Query(nativeQuery = true, value = "SELECT * FROM rule_node r WHERE r.type = :ruleType " + - " AND configuration_version < :version " + + " AND r.configuration_version < :version " + " AND (:searchText IS NULL OR r.configuration ILIKE CONCAT('%', :searchText, '%'))") Page findAllRuleNodesByTypeAndVersionLessThan(@Param("ruleType") String ruleType, @Param("version") int version, @@ -55,6 +55,11 @@ public interface RuleNodeRepository extends JpaRepository @Param("version") int version, Pageable pageable); + @Query("SELECT r.id FROM RuleNodeEntity r WHERE r.type = :ruleType AND r.configurationVersion = :version") + Page findAllRuleNodeIdsByTypeAndVersion(@Param("ruleType") String ruleType, + @Param("version") int version, + Pageable pageable); + List findRuleNodesByRuleChainIdAndExternalIdIn(UUID ruleChainId, List externalIds); @Transactional diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDaoTest.java index 2f1c73da01..8cb5974e7d 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDaoTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDaoTest.java @@ -144,6 +144,25 @@ public class JpaRuleNodeDaoTest extends AbstractJpaDaoTest { assertEquals(10, ruleNodeIds.getData().size()); } + @Test + public void testFindRuleNodeIdsByTypeAndVersion() { + PageData ruleNodeIds = ruleNodeDao.findAllRuleNodeIdsByTypeAndVersion( "A", 0, new PageLink(10, 0, PREFIX_FOR_RULE_NODE_NAME)); + assertEquals(20, ruleNodeIds.getTotalElements()); + assertEquals(2, ruleNodeIds.getTotalPages()); + assertEquals(10, ruleNodeIds.getData().size()); + + ruleNodeIds = ruleNodeDao.findAllRuleNodeIdsByTypeAndVersion( "A", 0, new PageLink(10, 0)); + assertEquals(20, ruleNodeIds.getTotalElements()); + assertEquals(2, ruleNodeIds.getTotalPages()); + assertEquals(10, ruleNodeIds.getData().size()); + + // test - search text ignored + ruleNodeIds = ruleNodeDao.findAllRuleNodeIdsByTypeAndVersion( "A", 0, new PageLink(10, 0, StringUtils.randomAlphabetic(5))); + assertEquals(20, ruleNodeIds.getTotalElements()); + assertEquals(2, ruleNodeIds.getTotalPages()); + assertEquals(10, ruleNodeIds.getData().size()); + } + @Test public void testFindAllRuleNodeByIds() { var fromUUIDs = ruleNodeIds.stream().map(RuleNodeId::new).collect(Collectors.toList()); From 7c42bb72138536c4f78351c1c992ec0d1780d8ea Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Fri, 29 Dec 2023 10:55:38 +0200 Subject: [PATCH 06/11] Change for loop to iterate forwards --- .../server/service/install/update/DefaultDataUpdateService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java index d0f45ac06c..a973ec023b 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java @@ -314,7 +314,7 @@ public class DefaultDataUpdateService implements DataUpdateService { private List getRuleNodesIdsWithTypeAndVersionLessThan(String type, int toVersion) { var ruleNodeIds = new ArrayList(); - for (int v = toVersion - 1; v >= 0; v--) { + for (int v = 0; v < toVersion; v++) { ruleNodeIds.addAll(getRuleNodesIdsWithTypeAndVersion(type, v)); } return ruleNodeIds; From 5e471364de2ba0263d62a144376174bcd0176b26 Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Tue, 2 Jan 2024 13:58:16 +0200 Subject: [PATCH 07/11] changed fixedPool to newWorkStealingPool --- .../java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java index 1c70302d00..3ccbc229a0 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java @@ -26,6 +26,8 @@ import com.amazonaws.services.sqs.model.CreateQueueRequest; import com.amazonaws.services.sqs.model.GetQueueUrlResult; import lombok.Getter; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.common.util.ExecutorProvider; +import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.util.PropertyUtils; @@ -55,7 +57,7 @@ public class TbAwsSqsAdmin implements TbQueueAdmin { AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); credentialsProvider = new AWSStaticCredentialsProvider(awsCredentials); } - producerExecutor = Executors.newFixedThreadPool(sqsSettings.getThreadPoolSize(), ThingsBoardThreadFactory.forName("aws-sqs-queue-executor")); + producerExecutor = ThingsBoardExecutors.newWorkStealingPool(sqsSettings.getThreadPoolSize(), "aws-sqs-queue-executor"); sqsClient = AmazonSQSClientBuilder.standard() .withCredentials(credentialsProvider) From 4a26fad0a53abb37cf6eaacd42a14d014619fb9a Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 2 Jan 2024 15:31:30 +0200 Subject: [PATCH 08/11] Revert "replaced slow query with multi-call of indexed query" This reverts commit abfcf32c54e2705171d93abbb563fbdec97def5d and 7c42bb72138536c4f78351c1c992ec0d1780d8ea. --- .../update/DefaultDataUpdateService.java | 10 +--------- .../server/dao/rule/RuleChainService.java | 4 ---- .../server/dao/rule/BaseRuleChainService.java | 10 ---------- .../server/dao/rule/RuleNodeDao.java | 2 -- .../server/dao/service/Validator.java | 15 +-------------- .../server/dao/sql/rule/JpaRuleNodeDao.java | 10 ---------- .../dao/sql/rule/RuleNodeRepository.java | 7 +------ .../dao/sql/rule/JpaRuleNodeDaoTest.java | 19 ------------------- 8 files changed, 3 insertions(+), 74 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java index a973ec023b..7d08fdca00 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java @@ -313,17 +313,9 @@ public class DefaultDataUpdateService implements DataUpdateService { } private List getRuleNodesIdsWithTypeAndVersionLessThan(String type, int toVersion) { - var ruleNodeIds = new ArrayList(); - for (int v = 0; v < toVersion; v++) { - ruleNodeIds.addAll(getRuleNodesIdsWithTypeAndVersion(type, v)); - } - return ruleNodeIds; - } - - private List getRuleNodesIdsWithTypeAndVersion(String type, int version) { var ruleNodeIds = new ArrayList(); new PageDataIterable<>(pageLink -> - ruleChainService.findAllRuleNodeIdsByTypeAndVersion(type, version, pageLink), DEFAULT_PAGE_SIZE + ruleChainService.findAllRuleNodeIdsByTypeAndVersionLessThan(type, toVersion, pageLink), DEFAULT_PAGE_SIZE ).forEach(ruleNodeIds::add); return ruleNodeIds; } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java index 03886bfecc..532da4ac00 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java @@ -100,14 +100,10 @@ public interface RuleChainService extends EntityDaoService { PageData findAllRuleNodesByType(String type, PageLink pageLink); - @Deprecated(forRemoval = true, since = "3.6.3") PageData findAllRuleNodesByTypeAndVersionLessThan(String type, int version, PageLink pageLink); - @Deprecated(forRemoval = true, since = "3.6.3") PageData findAllRuleNodeIdsByTypeAndVersionLessThan(String type, int version, PageLink pageLink); - PageData findAllRuleNodeIdsByTypeAndVersion(String type, int version, PageLink pageLink); - List findAllRuleNodesByIds(List ruleNodeIds); RuleNode saveRuleNode(TenantId tenantId, RuleNode ruleNode); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 5a6d0c1d74..6257bb0420 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -81,7 +81,6 @@ import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateIds; import static org.thingsboard.server.dao.service.Validator.validatePageLink; import static org.thingsboard.server.dao.service.Validator.validatePositiveNumber; -import static org.thingsboard.server.dao.service.Validator.validateNonNegativeNumber; import static org.thingsboard.server.dao.service.Validator.validateString; /** @@ -738,15 +737,6 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC return ruleNodeDao.findAllRuleNodeIdsByTypeAndVersionLessThan(type, version, pageLink); } - @Override - public PageData findAllRuleNodeIdsByTypeAndVersion(String type, int version, PageLink pageLink) { - log.trace("Executing findAllRuleNodeIdsByTypeAndVersion, type {}, pageLink {}, version {}", type, pageLink, version); - validateString(type, "Incorrect type of the rule node"); - validateNonNegativeNumber(version, "Incorrect version. Version should be non-negative!"); - validatePageLink(pageLink); - return ruleNodeDao.findAllRuleNodeIdsByTypeAndVersion(type, version, pageLink); - } - @Override public List findAllRuleNodesByIds(List ruleNodeIds) { log.trace("Executing findAllRuleNodesByIds, ruleNodeIds {}", ruleNodeIds); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeDao.java b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeDao.java index 4927f2d68d..b6f2dd097e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeDao.java @@ -38,8 +38,6 @@ public interface RuleNodeDao extends Dao { PageData findAllRuleNodeIdsByTypeAndVersionLessThan(String type, int version, PageLink pageLink); - PageData findAllRuleNodeIdsByTypeAndVersion(String type, int version, PageLink pageLink); - List findAllRuleNodeByIds(List ruleNodeIds); List findByExternalIds(RuleChainId ruleChainId, List externalIds); diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java b/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java index a712f7c775..d8c4ffddc9 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java @@ -61,7 +61,7 @@ public class Validator { /** - * This method validate long value. If value isn't positive than throw + * This method validate long value. If value isn't possitive than throw * IncorrectParameterException exception * * @param val the val @@ -73,19 +73,6 @@ public class Validator { } } - /** - * This method validate long value. If value is negative than throw - * IncorrectParameterException exception - * - * @param val the val - * @param errorMessage the error message for exception - */ - public static void validateNonNegativeNumber(long val, String errorMessage) { - if (val < 0) { - throw new IncorrectParameterException(errorMessage); - } - } - /** * This method validate UUID id. If id is null than throw * IncorrectParameterException exception diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java index 655a03caf3..de5ff109d3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java @@ -89,16 +89,6 @@ public class JpaRuleNodeDao extends JpaAbstractDao imp .mapData(RuleNodeId::new); } - @Override - public PageData findAllRuleNodeIdsByTypeAndVersion(String type, int version, PageLink pageLink) { - return DaoUtil.pageToPageData(ruleNodeRepository - .findAllRuleNodeIdsByTypeAndVersion( - type, - version, - DaoUtil.toPageable(pageLink))) - .mapData(RuleNodeId::new); - } - @Override public List findAllRuleNodeByIds(List ruleNodeIds) { return DaoUtil.convertDataList(ruleNodeRepository.findAllById( diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java index eaccf39072..f02072b356 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java @@ -43,7 +43,7 @@ public interface RuleNodeRepository extends JpaRepository Pageable pageable); @Query(nativeQuery = true, value = "SELECT * FROM rule_node r WHERE r.type = :ruleType " + - " AND r.configuration_version < :version " + + " AND configuration_version < :version " + " AND (:searchText IS NULL OR r.configuration ILIKE CONCAT('%', :searchText, '%'))") Page findAllRuleNodesByTypeAndVersionLessThan(@Param("ruleType") String ruleType, @Param("version") int version, @@ -55,11 +55,6 @@ public interface RuleNodeRepository extends JpaRepository @Param("version") int version, Pageable pageable); - @Query("SELECT r.id FROM RuleNodeEntity r WHERE r.type = :ruleType AND r.configurationVersion = :version") - Page findAllRuleNodeIdsByTypeAndVersion(@Param("ruleType") String ruleType, - @Param("version") int version, - Pageable pageable); - List findRuleNodesByRuleChainIdAndExternalIdIn(UUID ruleChainId, List externalIds); @Transactional diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDaoTest.java index 8cb5974e7d..2f1c73da01 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDaoTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDaoTest.java @@ -144,25 +144,6 @@ public class JpaRuleNodeDaoTest extends AbstractJpaDaoTest { assertEquals(10, ruleNodeIds.getData().size()); } - @Test - public void testFindRuleNodeIdsByTypeAndVersion() { - PageData ruleNodeIds = ruleNodeDao.findAllRuleNodeIdsByTypeAndVersion( "A", 0, new PageLink(10, 0, PREFIX_FOR_RULE_NODE_NAME)); - assertEquals(20, ruleNodeIds.getTotalElements()); - assertEquals(2, ruleNodeIds.getTotalPages()); - assertEquals(10, ruleNodeIds.getData().size()); - - ruleNodeIds = ruleNodeDao.findAllRuleNodeIdsByTypeAndVersion( "A", 0, new PageLink(10, 0)); - assertEquals(20, ruleNodeIds.getTotalElements()); - assertEquals(2, ruleNodeIds.getTotalPages()); - assertEquals(10, ruleNodeIds.getData().size()); - - // test - search text ignored - ruleNodeIds = ruleNodeDao.findAllRuleNodeIdsByTypeAndVersion( "A", 0, new PageLink(10, 0, StringUtils.randomAlphabetic(5))); - assertEquals(20, ruleNodeIds.getTotalElements()); - assertEquals(2, ruleNodeIds.getTotalPages()); - assertEquals(10, ruleNodeIds.getData().size()); - } - @Test public void testFindAllRuleNodeByIds() { var fromUUIDs = ruleNodeIds.stream().map(RuleNodeId::new).collect(Collectors.toList()); From 360e32684ed8458de375a953998dd258100c92ea Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 2 Jan 2024 15:34:32 +0200 Subject: [PATCH 09/11] added new index for rule nodes ids search query and removed old one --- .../main/data/upgrade/3.6.2/schema_update.sql | 22 +++++++++++++++++++ .../install/SqlDatabaseUpgradeService.java | 3 +++ .../server/dao/rule/RuleChainService.java | 1 + .../server/dao/service/Validator.java | 2 +- .../dao/sql/rule/RuleNodeRepository.java | 2 +- .../resources/sql/schema-entities-idx.sql | 2 +- 6 files changed, 29 insertions(+), 3 deletions(-) create mode 100644 application/src/main/data/upgrade/3.6.2/schema_update.sql diff --git a/application/src/main/data/upgrade/3.6.2/schema_update.sql b/application/src/main/data/upgrade/3.6.2/schema_update.sql new file mode 100644 index 0000000000..fd64136852 --- /dev/null +++ b/application/src/main/data/upgrade/3.6.2/schema_update.sql @@ -0,0 +1,22 @@ +-- +-- Copyright © 2016-2023 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. +-- + +-- RULE NODE INDEXES UPDATE START + +DROP INDEX IF EXISTS idx_rule_node_type_configuration_version; +CREATE INDEX IF NOT EXISTS idx_rule_node_id_type_configuration_version ON rule_node(id, type, configuration_version); + +-- RULE NODE INDEXES UPDATE END diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index c2444390f6..4a96d0d2c3 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -772,6 +772,9 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService } }); break; + case "3.6.2": + updateSchema("3.6.2", 3006002, "3.6.3", 3006003, null); + break; default: throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion); } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java index 532da4ac00..a8a60bc565 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java @@ -100,6 +100,7 @@ public interface RuleChainService extends EntityDaoService { PageData findAllRuleNodesByType(String type, PageLink pageLink); + @Deprecated(forRemoval = true, since = "3.6.3") PageData findAllRuleNodesByTypeAndVersionLessThan(String type, int version, PageLink pageLink); PageData findAllRuleNodeIdsByTypeAndVersionLessThan(String type, int version, PageLink pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java b/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java index d8c4ffddc9..3c238305bd 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java @@ -61,7 +61,7 @@ public class Validator { /** - * This method validate long value. If value isn't possitive than throw + * This method validate long value. If value isn't positive than throw * IncorrectParameterException exception * * @param val the val diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java index f02072b356..de61422527 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java @@ -43,7 +43,7 @@ public interface RuleNodeRepository extends JpaRepository Pageable pageable); @Query(nativeQuery = true, value = "SELECT * FROM rule_node r WHERE r.type = :ruleType " + - " AND configuration_version < :version " + + " AND r.configuration_version < :version " + " AND (:searchText IS NULL OR r.configuration ILIKE CONCAT('%', :searchText, '%'))") Page findAllRuleNodesByTypeAndVersionLessThan(@Param("ruleType") String ruleType, @Param("version") int version, diff --git a/dao/src/main/resources/sql/schema-entities-idx.sql b/dao/src/main/resources/sql/schema-entities-idx.sql index df95a50d47..10678509b4 100644 --- a/dao/src/main/resources/sql/schema-entities-idx.sql +++ b/dao/src/main/resources/sql/schema-entities-idx.sql @@ -93,7 +93,7 @@ CREATE INDEX IF NOT EXISTS idx_rule_node_external_id ON rule_node(rule_chain_id, CREATE INDEX IF NOT EXISTS idx_rule_node_type ON rule_node(type); -CREATE INDEX IF NOT EXISTS idx_rule_node_type_configuration_version ON rule_node(type, configuration_version); +CREATE INDEX IF NOT EXISTS idx_rule_node_id_type_configuration_version ON rule_node(id, type, configuration_version); CREATE INDEX IF NOT EXISTS idx_api_usage_state_entity_id ON api_usage_state(entity_id); From f14e6cdfcae7caed69a70b7c362d6fc2c409e4f6 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 2 Jan 2024 15:56:37 +0200 Subject: [PATCH 10/11] added missing switch case to ThingsboardInstallService --- .../thingsboard/server/install/ThingsboardInstallService.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index 5ab3b7b28e..8b14d2d8b7 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java @@ -276,6 +276,9 @@ public class ThingsboardInstallService { } else { log.info("Skipping images migration. Run the upgrade with fromVersion as '3.6.2-images' to migrate"); } + case "3.6.2": + log.info("Upgrading ThingsBoard from version 3.6.2 to 3.6.3 ..."); + databaseEntitiesUpgradeService.upgradeDatabase("3.6.2"); //TODO DON'T FORGET to update switch statement in the CacheCleanupService if you need to clear the cache break; default: From 14b3dd4b57edaba259ec84494d4bf10011cfae8d Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Wed, 3 Jan 2024 10:36:08 +0200 Subject: [PATCH 11/11] updated columns order in index --- application/src/main/data/upgrade/3.6.2/schema_update.sql | 3 ++- dao/src/main/resources/sql/schema-entities-idx.sql | 4 +--- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/application/src/main/data/upgrade/3.6.2/schema_update.sql b/application/src/main/data/upgrade/3.6.2/schema_update.sql index fd64136852..4e62209ea8 100644 --- a/application/src/main/data/upgrade/3.6.2/schema_update.sql +++ b/application/src/main/data/upgrade/3.6.2/schema_update.sql @@ -16,7 +16,8 @@ -- RULE NODE INDEXES UPDATE START +DROP INDEX IF EXISTS idx_rule_node_type; DROP INDEX IF EXISTS idx_rule_node_type_configuration_version; -CREATE INDEX IF NOT EXISTS idx_rule_node_id_type_configuration_version ON rule_node(id, type, configuration_version); +CREATE INDEX IF NOT EXISTS idx_rule_node_type_id_configuration_version ON rule_node(type, id, configuration_version); -- RULE NODE INDEXES UPDATE END diff --git a/dao/src/main/resources/sql/schema-entities-idx.sql b/dao/src/main/resources/sql/schema-entities-idx.sql index 10678509b4..6a95d9e346 100644 --- a/dao/src/main/resources/sql/schema-entities-idx.sql +++ b/dao/src/main/resources/sql/schema-entities-idx.sql @@ -91,9 +91,7 @@ CREATE INDEX IF NOT EXISTS idx_widgets_bundle_external_id ON widgets_bundle(tena CREATE INDEX IF NOT EXISTS idx_rule_node_external_id ON rule_node(rule_chain_id, external_id); -CREATE INDEX IF NOT EXISTS idx_rule_node_type ON rule_node(type); - -CREATE INDEX IF NOT EXISTS idx_rule_node_id_type_configuration_version ON rule_node(id, type, configuration_version); +CREATE INDEX IF NOT EXISTS idx_rule_node_type_id_configuration_version ON rule_node(type, id, configuration_version); CREATE INDEX IF NOT EXISTS idx_api_usage_state_entity_id ON api_usage_state(entity_id);