diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java index f108de4a43..3355c55d7e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java @@ -67,6 +67,7 @@ public class TbServiceBusAdmin implements TbQueueAdmin { try { QueueDescription queueDescription = new QueueDescription(topic); + queueDescription.setRequiresDuplicateDetection(false); setQueueConfigs(queueDescription); client.createQueue(queueDescription); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java index 4db5ade728..734c79d927 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java @@ -58,7 +58,7 @@ public class TbServiceBusConsumerTemplate extends Abstract private final Gson gson = new Gson(); private Set receivers; - private Map> pendingMessages = new ConcurrentHashMap<>(); + private final Map> pendingMessages = new ConcurrentHashMap<>(); private volatile int messagesPerQueue; public TbServiceBusConsumerTemplate(TbQueueAdmin admin, TbServiceBusSettings serviceBusSettings, String topic, TbQueueMsgDecoder decoder) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusProducerTemplate.java index 5d9a931378..3c5228198d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusProducerTemplate.java @@ -33,6 +33,7 @@ import org.thingsboard.server.queue.common.DefaultTbQueueMsg; import java.util.HashMap; import java.util.Map; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -42,14 +43,14 @@ public class TbServiceBusProducerTemplate implements TbQue private final Gson gson = new Gson(); private final TbQueueAdmin admin; private final TbServiceBusSettings serviceBusSettings; - private final Map clients = new HashMap<>(); - private ExecutorService executorService; + private final Map clients = new ConcurrentHashMap<>(); + private final ExecutorService executorService; public TbServiceBusProducerTemplate(TbQueueAdmin admin, TbServiceBusSettings serviceBusSettings, String defaultTopic) { this.admin = admin; this.defaultTopic = defaultTopic; this.serviceBusSettings = serviceBusSettings; - executorService = Executors.newSingleThreadExecutor(); + executorService = Executors.newCachedThreadPool(); } @Override diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java index 75635de7a4..1b5619fd26 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java @@ -21,7 +21,6 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; -import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueMsg; import org.thingsboard.server.queue.common.AbstractTbQueueConsumerTemplate; @@ -32,7 +31,6 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Properties; -import java.util.stream.Collectors; /** * Created by ashvayka on 24.09.18. @@ -69,7 +67,7 @@ public class TbKafkaConsumerTemplate extends AbstractTbQue } @Override - protected void doSubscribe( List topicNames) { + protected void doSubscribe(List topicNames) { topicNames.forEach(admin::createTopicIfNotExists); consumer.subscribe(topicNames); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java index 1241370d05..d0a514ffd3 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.queue.pubsub; +import com.google.api.gax.rpc.AlreadyExistsException; import com.google.cloud.pubsub.v1.SubscriptionAdminClient; import com.google.cloud.pubsub.v1.SubscriptionAdminSettings; import com.google.cloud.pubsub.v1.TopicAdminClient; @@ -24,9 +25,9 @@ import com.google.pubsub.v1.ListSubscriptionsRequest; import com.google.pubsub.v1.ListTopicsRequest; import com.google.pubsub.v1.ProjectName; import com.google.pubsub.v1.ProjectSubscriptionName; -import com.google.pubsub.v1.ProjectTopicName; import com.google.pubsub.v1.Subscription; import com.google.pubsub.v1.Topic; +import com.google.pubsub.v1.TopicName; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.queue.TbQueueAdmin; @@ -103,7 +104,10 @@ public class TbPubSubAdmin implements TbQueueAdmin { @Override public void createTopicIfNotExists(String partition) { - ProjectTopicName topicName = ProjectTopicName.of(pubSubSettings.getProjectId(), partition); + TopicName topicName = TopicName.newBuilder() + .setTopic(partition) + .setProject(pubSubSettings.getProjectId()) + .build(); if (topicSet.contains(topicName.toString())) { createSubscriptionIfNotExists(partition, topicName); @@ -121,13 +125,18 @@ public class TbPubSubAdmin implements TbQueueAdmin { } } - topicAdminClient.createTopic(topicName); - topicSet.add(topicName.toString()); - log.info("Created new topic: [{}]", topicName.toString()); + try { + topicAdminClient.createTopic(topicName); + log.info("Created new topic: [{}]", topicName.toString()); + } catch (AlreadyExistsException e) { + log.info("[{}] Topic already exist.", topicName.toString()); + } finally { + topicSet.add(topicName.toString()); + } createSubscriptionIfNotExists(partition, topicName); } - private void createSubscriptionIfNotExists(String partition, ProjectTopicName topicName) { + private void createSubscriptionIfNotExists(String partition, TopicName topicName) { ProjectSubscriptionName subscriptionName = ProjectSubscriptionName.of(pubSubSettings.getProjectId(), partition); @@ -153,9 +162,14 @@ public class TbPubSubAdmin implements TbQueueAdmin { setAckDeadline(subscriptionBuilder); setMessageRetention(subscriptionBuilder); - subscriptionAdminClient.createSubscription(subscriptionBuilder.build()); - subscriptionSet.add(subscriptionName.toString()); - log.info("Created new subscription: [{}]", subscriptionName.toString()); + try { + subscriptionAdminClient.createSubscription(subscriptionBuilder.build()); + log.info("Created new subscription: [{}]", subscriptionName.toString()); + } catch (AlreadyExistsException e) { + log.info("[{}] Subscription already exist.", subscriptionName.toString()); + } finally { + subscriptionSet.add(subscriptionName.toString()); + } } private void setAckDeadline(Subscription.Builder builder) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java index 7302d19ff7..b5b6126cd5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.queue.pubsub; -import com.amazonaws.services.sqs.model.Message; import com.google.api.core.ApiFuture; import com.google.api.core.ApiFutures; import com.google.cloud.pubsub.v1.stub.GrpcSubscriberStub; @@ -31,13 +30,10 @@ import com.google.pubsub.v1.PullResponse; import com.google.pubsub.v1.ReceivedMessage; import lombok.extern.slf4j.Slf4j; import org.springframework.util.CollectionUtils; -import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.queue.TbQueueAdmin; -import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueMsg; import org.thingsboard.server.queue.TbQueueMsgDecoder; import org.thingsboard.server.queue.common.AbstractParallelTbQueueConsumerTemplate; -import org.thingsboard.server.queue.common.AbstractTbQueueConsumerTemplate; import org.thingsboard.server.queue.common.DefaultTbQueueMsg; import java.io.IOException; @@ -49,9 +45,6 @@ import java.util.Objects; import java.util.Set; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.ExecutionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @Slf4j @@ -136,7 +129,7 @@ public class TbPubSubConsumerTemplate extends AbstractPara PullRequest pullRequest = PullRequest.newBuilder() .setMaxMessages(messagesPerTopic) - .setReturnImmediately(false) // return immediately if messages are not available +// .setReturnImmediately(false) // return immediately if messages are not available .setSubscription(subscriptionName) .build(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubProducerTemplate.java index 2cd2e1054e..ea713827e5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubProducerTemplate.java @@ -49,7 +49,7 @@ public class TbPubSubProducerTemplate implements TbQueuePr private final Map publisherMap = new ConcurrentHashMap<>(); - private ExecutorService pubExecutor = Executors.newCachedThreadPool(); + private final ExecutorService pubExecutor = Executors.newCachedThreadPool(); public TbPubSubProducerTemplate(TbQueueAdmin admin, TbPubSubSettings pubSubSettings, String defaultTopic) { this.defaultTopic = defaultTopic; @@ -124,8 +124,8 @@ public class TbPubSubProducerTemplate implements TbQueuePr publisherMap.put(topic, publisher); return publisher; } catch (IOException e) { - log.error("Failed to create topic [{}].", topic, e); - throw new RuntimeException("Failed to create topic.", e); + log.error("Failed to create Publisher for the topic [{}].", topic, e); + throw new RuntimeException("Failed to create Publisher for the topic.", e); } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java index 3bef6deb84..676ee5354f 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java @@ -27,13 +27,11 @@ import java.util.concurrent.TimeoutException; @Slf4j public class TbRabbitMqAdmin implements TbQueueAdmin { - private final TbRabbitMqSettings rabbitMqSettings; private final Channel channel; private final Connection connection; private final Map arguments; public TbRabbitMqAdmin(TbRabbitMqSettings rabbitMqSettings, Map arguments) { - this.rabbitMqSettings = rabbitMqSettings; this.arguments = arguments; try { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqProducerTemplate.java index 91b46213a5..a58b817f78 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqProducerTemplate.java @@ -30,6 +30,8 @@ import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.DefaultTbQueueMsg; import java.io.IOException; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Executors; import java.util.concurrent.TimeoutException; @@ -39,10 +41,12 @@ public class TbRabbitMqProducerTemplate implements TbQueue private final Gson gson = new Gson(); private final TbQueueAdmin admin; private final TbRabbitMqSettings rabbitMqSettings; - private ListeningExecutorService producerExecutor; + private final ListeningExecutorService producerExecutor; private final Channel channel; private final Connection connection; + private final Set topics = ConcurrentHashMap.newKeySet(); + public TbRabbitMqProducerTemplate(TbQueueAdmin admin, TbRabbitMqSettings rabbitMqSettings, String defaultTopic) { this.admin = admin; this.defaultTopic = defaultTopic; @@ -75,6 +79,7 @@ public class TbRabbitMqProducerTemplate implements TbQueue @Override public void send(TopicPartitionInfo tpi, T msg, TbQueueCallback callback) { + createTopicIfNotExist(tpi); AMQP.BasicProperties properties = new AMQP.BasicProperties(); try { channel.basicPublish(rabbitMqSettings.getExchangeName(), tpi.getFullTopicName(), properties, gson.toJson(new DefaultTbQueueMsg(msg)).getBytes()); @@ -110,4 +115,11 @@ public class TbRabbitMqProducerTemplate implements TbQueue } } + private void createTopicIfNotExist(TopicPartitionInfo tpi) { + if (topics.contains(tpi)) { + return; + } + admin.createTopicIfNotExists(tpi.getFullTopicName()); + topics.add(tpi); + } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsConsumerTemplate.java index 317dd93902..f4f279a02b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsConsumerTemplate.java @@ -84,10 +84,6 @@ public class TbAwsSqsConsumerTemplate extends AbstractPara @Override protected List doPoll(long durationInMillis) { - if (!pendingMessages.isEmpty()) { - log.warn("Present {} non committed messages.", pendingMessages.size()); - return Collections.emptyList(); - } int duration = (int) TimeUnit.MILLISECONDS.toSeconds(durationInMillis); List>> futureList = queueUrls .stream() @@ -145,7 +141,6 @@ public class TbAwsSqsConsumerTemplate extends AbstractPara ReceiveMessageRequest request = new ReceiveMessageRequest(); request .withWaitTimeSeconds(waitTimeSeconds) - .withMessageAttributeNames("headers") .withQueueUrl(url) .withMaxNumberOfMessages(MAX_NUM_MSGS); return sqsClient.receiveMessage(request).getMessages(); 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 2d85539184..6110d08c5e 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 @@ -37,6 +37,7 @@ import org.thingsboard.server.queue.TbQueueProducer; 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; @@ -80,7 +81,9 @@ public class TbAwsSqsProducerTemplate implements TbQueuePr sendMsgRequest.withQueueUrl(getQueueUrl(tpi.getFullTopicName())); sendMsgRequest.withMessageBody(gson.toJson(new DefaultTbQueueMsg(msg))); - sendMsgRequest.withMessageGroupId(msg.getKey().toString()); + sendMsgRequest.withMessageGroupId(tpi.getTopic()); + sendMsgRequest.withMessageDeduplicationId(UUID.randomUUID().toString()); + ListenableFuture future = producerExecutor.submit(() -> sqsClient.sendMessage(sendMsgRequest)); Futures.addCallback(future, new FutureCallback() { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsQueueAttributes.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsQueueAttributes.java index 70a9587bef..c6cbbfd256 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsQueueAttributes.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsQueueAttributes.java @@ -55,7 +55,6 @@ public class TbAwsSqsQueueAttributes { @PostConstruct private void init() { defaultAttributes.put(QueueAttributeName.FifoQueue.toString(), "true"); - defaultAttributes.put(QueueAttributeName.ContentBasedDeduplication.toString(), "true"); coreAttributes = getConfigs(coreProperties); ruleEngineAttributes = getConfigs(ruleEngineProperties); diff --git a/docker/tb-js-executor.env b/docker/tb-js-executor.env index 0b64f43b00..b66073ea44 100644 --- a/docker/tb-js-executor.env +++ b/docker/tb-js-executor.env @@ -1,4 +1,4 @@ - +TB_QUEUE_TYPE=kafka REMOTE_JS_EVAL_REQUEST_TOPIC=js_eval.requests TB_KAFKA_SERVERS=kafka:9092 LOGGER_LEVEL=info diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index b290719739..c573274801 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -14,7 +14,7 @@ # limitations under the License. # -service-type: "TB_SERVICE_TYPE" #kafka (Apache Kafka) or aws-sqs (AWS SQS) or pubsub (PubSub) or service-bus (Azure Service Bus) or rabbitmq (RabbitMQ) +queue_type: "TB_QUEUE_TYPE" #kafka (Apache Kafka) or aws-sqs (AWS SQS) or pubsub (PubSub) or service-bus (Azure Service Bus) or rabbitmq (RabbitMQ) request_topic: "REMOTE_JS_EVAL_REQUEST_TOPIC" js: @@ -25,18 +25,18 @@ kafka: # Kafka Bootstrap Servers servers: "TB_KAFKA_SERVERS" replication_factor: "TB_QUEUE_KAFKA_REPLICATION_FACTOR" - topic-properties: "TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES" + topic_properties: "TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES" pubsub: project_id: "TB_QUEUE_PUBSUB_PROJECT_ID" service_account: "TB_QUEUE_PUBSUB_SERVICE_ACCOUNT" - queue-properties: "TB_QUEUE_PUBSUB_JE_QUEUE_PROPERTIES" + queue_properties: "TB_QUEUE_PUBSUB_JE_QUEUE_PROPERTIES" aws_sqs: access_key_id: "TB_QUEUE_AWS_SQS_ACCESS_KEY_ID" secret_access_key: "TB_QUEUE_AWS_SQS_SECRET_ACCESS_KEY" region: "TB_QUEUE_AWS_SQS_REGION" - queue-properties: "TB_QUEUE_AWS_SQS_JE_QUEUE_PROPERTIES" + queue_properties: "TB_QUEUE_AWS_SQS_JE_QUEUE_PROPERTIES" rabbitmq: host: "TB_QUEUE_RABBIT_MQ_HOST" @@ -44,14 +44,14 @@ rabbitmq: virtual_host: "TB_QUEUE_RABBIT_MQ_VIRTUAL_HOST" username: "TB_QUEUE_RABBIT_MQ_USERNAME" password: "TB_QUEUE_RABBIT_MQ_PASSWORD" - queue-properties: "TB_QUEUE_RABBIT_MQ_JE_QUEUE_PROPERTIES" + queue_properties: "TB_QUEUE_RABBIT_MQ_JE_QUEUE_PROPERTIES" service_bus: namespace_name: "TB_QUEUE_SERVICE_BUS_NAMESPACE_NAME" sas_key_name: "TB_QUEUE_SERVICE_BUS_SAS_KEY_NAME" sas_key: "TB_QUEUE_SERVICE_BUS_SAS_KEY" max_messages: "TB_QUEUE_SERVICE_BUS_MAX_MESSAGES" - queue-properties: "TB_QUEUE_SERVICE_BUS_JE_QUEUE_PROPERTIES" + queue_properties: "TB_QUEUE_SERVICE_BUS_JE_QUEUE_PROPERTIES" logger: level: "LOGGER_LEVEL" diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index 3155b051dc..f42b74745f 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -14,7 +14,7 @@ # limitations under the License. # -service-type: "kafka" +queue_type: "kafka" request_topic: "js_eval.requests" js: @@ -25,13 +25,13 @@ kafka: # Kafka Bootstrap Servers servers: "localhost:9092" replication_factor: "1" - topic-properties: "retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600" + topic_properties: "retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600" pubsub: - queue-properties: "ackDeadlineInSec:30;messageRetentionInSec:604800" + queue_properties: "ackDeadlineInSec:30;messageRetentionInSec:604800" aws_sqs: - queue-properties: "VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800" + queue_properties: "VisibilityTimeout:30;MaximumMessageSize:262144;MessageRetentionPeriod:604800" rabbitmq: host: "localhost" @@ -39,10 +39,10 @@ rabbitmq: virtual_host: "/" username: "admin" password: "password" - queue-properties: "x-max-length-bytes:1048576000;x-message-ttl:604800000" + queue_properties: "x-max-length-bytes:1048576000;x-message-ttl:604800000" service_bus: - queue-properties: "lockDurationInSec:30;maxSizeInMb:1024;messageTimeToLiveInSec:604800" + queue_properties: "lockDurationInSec:30;maxSizeInMb:1024;messageTimeToLiveInSec:604800" logger: level: "info" diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index 3dadac4f84..60a107061a 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -22,6 +22,7 @@ "azure-sb": "^0.11.1", "long": "^4.0.0", "uuid-parse": "^1.0.0", + "uuid-random": "^1.3.0", "winston": "^3.0.0", "winston-daily-rotate-file": "^3.2.1" }, diff --git a/msa/js-executor/queue/awsSqsTemplate.js b/msa/js-executor/queue/awsSqsTemplate.js index 0396824af1..a74d9d5b57 100644 --- a/msa/js-executor/queue/awsSqsTemplate.js +++ b/msa/js-executor/queue/awsSqsTemplate.js @@ -19,6 +19,7 @@ const config = require('config'), JsInvokeMessageProcessor = require('../api/jsInvokeMessageProcessor'), logger = require('../config/logger')._logger('awsSqsTemplate'); +const uuid = require('uuid-random'); const requestTopic = config.get('request_topic'); @@ -26,10 +27,10 @@ const accessKeyId = config.get('aws_sqs.access_key_id'); const secretAccessKey = config.get('aws_sqs.secret_access_key'); const region = config.get('aws_sqs.region'); const AWS = require('aws-sdk'); -const queueProperties = config.get('aws_sqs.queue-properties'); -const poolInterval = config.get('js.response_poll_interval'); +const queueProperties = config.get('aws_sqs.queue_properties'); +const pollInterval = config.get('js.response_poll_interval'); -let queueAttributes = {FifoQueue: 'true', ContentBasedDeduplication: 'true'}; +let queueAttributes = {FifoQueue: 'true'}; let sqsClient; let requestQueueURL; const queueUrls = new Map(); @@ -51,7 +52,12 @@ function AwsSqsProducer() { queueUrls.set(responseTopic, responseQueueUrl); } - let params = {MessageBody: msgBody, QueueUrl: responseQueueUrl, MessageGroupId: scriptId}; + let params = { + MessageBody: msgBody, + QueueUrl: responseQueueUrl, + MessageGroupId: 'js_eval', + MessageDeduplicationId: uuid() + }; return new Promise((resolve, reject) => { sqsClient.sendMessage(params, function (err, data) { @@ -74,11 +80,13 @@ function AwsSqsProducer() { const queues = await getQueues(); - queues.forEach(queueUrl => { - const delimiterPosition = queueUrl.lastIndexOf('/'); - const queueName = queueUrl.substring(delimiterPosition + 1); - queueUrls.set(queueName, queueUrl); - }) + if (queues) { + queues.forEach(queueUrl => { + const delimiterPosition = queueUrl.lastIndexOf('/'); + const queueName = queueUrl.substring(delimiterPosition + 1); + queueUrls.set(queueName, queueUrl); + }); + } parseQueueProperties(); @@ -95,6 +103,7 @@ function AwsSqsProducer() { WaitTimeSeconds: poolInterval / 1000 }; while (!stopped) { + let pollStartTs = new Date().getTime(); const messages = await new Promise((resolve, reject) => { sqsClient.receiveMessage(params, function (err, data) { if (err) { @@ -127,6 +136,11 @@ function AwsSqsProducer() { //do nothing } }); + } else { + let pollDuration = new Date().getTime() - pollStartTs; + if (pollDuration < pollInterval) { + await sleep(pollInterval - pollDuration); + } } } } catch (e) { @@ -175,6 +189,12 @@ function parseQueueProperties() { }); } +function sleep(ms) { + return new Promise((resolve) => { + setTimeout(resolve, ms); + }); +} + process.on('exit', () => { stopped = true; logger.info('Aws Sqs client stopped.'); diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 6a09338427..33637cee81 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -20,7 +20,7 @@ const config = require('config'), logger = require('../config/logger')._logger('kafkaTemplate'), KafkaJsWinstonLogCreator = require('../config/logger').KafkaJsWinstonLogCreator; const replicationFactor = config.get('kafka.replication_factor'); -const topicProperties = config.get('kafka.topic-properties'); +const topicProperties = config.get('kafka.topic_properties'); let kafkaClient; let kafkaAdmin; diff --git a/msa/js-executor/queue/pubSubTemplate.js b/msa/js-executor/queue/pubSubTemplate.js index 17e1b56e1d..8cb264de64 100644 --- a/msa/js-executor/queue/pubSubTemplate.js +++ b/msa/js-executor/queue/pubSubTemplate.js @@ -24,7 +24,7 @@ const {PubSub} = require('@google-cloud/pubsub'); const projectId = config.get('pubsub.project_id'); const credentials = JSON.parse(config.get('pubsub.service_account')); const requestTopic = config.get('request_topic'); -const queueProperties = config.get('pubsub.queue-properties'); +const queueProperties = config.get('pubsub.queue_properties'); let pubSubClient; @@ -98,23 +98,32 @@ function PubSubProducer() { async function createTopic(topic) { if (!topics.includes(topic)) { - await pubSubClient.createTopic(topic); + try { + await pubSubClient.createTopic(topic); + logger.info('Created new Pub/Sub topic: %s', topic); + } catch (e) { + logger.info('Pub/Sub topic already exists'); + } topics.push(topic); - logger.info('Created new Pub/Sub topic: %s', topic); } await createSubscription(topic) } async function createSubscription(topic) { if (!subscriptions.includes(topic)) { - await pubSubClient.createSubscription(topic, topic, { - topic: topic, - subscription: topic, - ackDeadlineSeconds: queueProps['ackDeadlineInSec'], - messageRetentionDuration: {seconds: queueProps['messageRetentionInSec']} - }); + try { + await pubSubClient.createSubscription(topic, topic, { + topic: topic, + subscription: topic, + ackDeadlineSeconds: queueProps['ackDeadlineInSec'], + messageRetentionDuration: {seconds: queueProps['messageRetentionInSec']} + }); + logger.info('Created new Pub/Sub subscription: %s', topic); + } catch (e) { + logger.info('Pub/Sub subscription already exists.'); + } + subscriptions.push(topic); - logger.info('Created new Pub/Sub subscription: %s', topic); } } diff --git a/msa/js-executor/queue/rabbitmqTemplate.js b/msa/js-executor/queue/rabbitmqTemplate.js index 1a2905c3a0..2aeb5c02db 100644 --- a/msa/js-executor/queue/rabbitmqTemplate.js +++ b/msa/js-executor/queue/rabbitmqTemplate.js @@ -26,23 +26,23 @@ const port = config.get('rabbitmq.port'); const vhost = config.get('rabbitmq.virtual_host'); const username = config.get('rabbitmq.username'); const password = config.get('rabbitmq.password'); -const queueProperties = config.get('rabbitmq.queue-properties'); -const poolInterval = config.get('js.response_poll_interval'); +const queueProperties = config.get('rabbitmq.queue_properties'); +const pollInterval = config.get('js.response_poll_interval'); const amqp = require('amqplib/callback_api'); -let queueParams = {durable: false, exclusive: false, autoDelete: false}; +let queueOptions = {durable: false, exclusive: false, autoDelete: false}; let connection; let channel; let stopped = false; -const responseTopics = []; +let queues = []; function RabbitMqProducer() { this.send = async (responseTopic, scriptId, rawResponse, headers) => { - if (!responseTopics.includes(responseTopic)) { + if (!queues.includes(responseTopic)) { await createQueue(responseTopic); - responseTopics.push(responseTopic); + queues.push(responseTopic); } let data = JSON.stringify( @@ -98,6 +98,7 @@ function RabbitMqProducer() { const messageProcessor = new JsInvokeMessageProcessor(new RabbitMqProducer()); while (!stopped) { + let pollStartTs = new Date().getTime(); let message = await new Promise((resolve, reject) => { channel.get(requestTopic, {}, function (err, msg) { if (err) { @@ -112,7 +113,10 @@ function RabbitMqProducer() { messageProcessor.onJsInvokeMessage(JSON.parse(message.content.toString('utf8'))); channel.ack(message); } else { - await sleep(poolInterval); + let pollDuration = new Date().getTime() - pollStartTs; + if (pollDuration < pollInterval) { + await sleep(pollInterval - pollDuration); + } } } } catch (e) { @@ -123,16 +127,18 @@ function RabbitMqProducer() { })(); function parseQueueProperties() { + let args = {}; const props = queueProperties.split(';'); props.forEach(p => { const delimiterPosition = p.indexOf(':'); - queueParams[p.substring(0, delimiterPosition)] = p.substring(delimiterPosition + 1); + args[p.substring(0, delimiterPosition)] = +p.substring(delimiterPosition + 1); }); + queueOptions['arguments'] = args; } -function createQueue(topic) { +async function createQueue(topic) { return new Promise((resolve, reject) => { - channel.assertQueue(topic, queueParams, function (err) { + channel.assertQueue(topic, queueOptions, function (err) { if (err) { reject(err); } else { diff --git a/msa/js-executor/queue/serviceBusTemplate.js b/msa/js-executor/queue/serviceBusTemplate.js index 034921afb7..99316930e3 100644 --- a/msa/js-executor/queue/serviceBusTemplate.js +++ b/msa/js-executor/queue/serviceBusTemplate.js @@ -26,7 +26,7 @@ const requestTopic = config.get('request_topic'); const namespaceName = config.get('service_bus.namespace_name'); const sasKeyName = config.get('service_bus.sas_key_name'); const sasKey = config.get('service_bus.sas_key'); -const queueProperties = config.get('service_bus.queue-properties'); +const queueProperties = config.get('service_bus.queue_properties'); let sbClient; let receiverClient; @@ -140,6 +140,7 @@ function parseQueueProperties() { properties[p.substring(0, delimiterPosition)] = p.substring(delimiterPosition + 1); }); queueOptions = { + DuplicateDetection: 'false', MaxSizeInMegabytes: properties['maxSizeInMb'], DefaultMessageTimeToLive: `PT${properties['messageTimeToLiveInSec']}S`, LockDuration: `PT${properties['lockDurationInSec']}S` diff --git a/msa/js-executor/server.js b/msa/js-executor/server.js index 58361016c4..e57ba62c11 100644 --- a/msa/js-executor/server.js +++ b/msa/js-executor/server.js @@ -16,7 +16,7 @@ const config = require('config'), logger = require('./config/logger')._logger('main'); -const serviceType = config.get('service-type'); +const serviceType = config.get('queue_type'); switch (serviceType) { case 'kafka': logger.info('Starting kafka template.'); diff --git a/pom.xml b/pom.xml index 494a43cb56..a1dbc2d8f5 100755 --- a/pom.xml +++ b/pom.xml @@ -95,7 +95,7 @@ 1.25 1.3.10 1.11.747 - 1.84.0 + 1.105.0 3.2.0 1.5.0 1.4.3