From 9ef3445b7723a48ad2b96d6cf3d1b3830e946635 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Thu, 30 Apr 2020 16:24:12 +0300 Subject: [PATCH 1/4] refactored --- .../config/custom-environment-variables.yml | 12 ++++++------ msa/js-executor/config/default.yml | 12 ++++++------ msa/js-executor/queue/awsSqsTemplate.js | 2 +- msa/js-executor/queue/kafkaTemplate.js | 2 +- msa/js-executor/queue/pubSubTemplate.js | 2 +- msa/js-executor/queue/rabbitmqTemplate.js | 2 +- msa/js-executor/queue/serviceBusTemplate.js | 2 +- msa/js-executor/server.js | 2 +- 8 files changed, 18 insertions(+), 18 deletions(-) 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/queue/awsSqsTemplate.js b/msa/js-executor/queue/awsSqsTemplate.js index 0396824af1..a5338c5e73 100644 --- a/msa/js-executor/queue/awsSqsTemplate.js +++ b/msa/js-executor/queue/awsSqsTemplate.js @@ -26,7 +26,7 @@ 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 queueProperties = config.get('aws_sqs.queue_properties'); const poolInterval = config.get('js.response_poll_interval'); let queueAttributes = {FifoQueue: 'true', ContentBasedDeduplication: 'true'}; diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 699ca7be00..e22ec8cd71 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..7d0b32ea34 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; diff --git a/msa/js-executor/queue/rabbitmqTemplate.js b/msa/js-executor/queue/rabbitmqTemplate.js index 1a2905c3a0..732206ff11 100644 --- a/msa/js-executor/queue/rabbitmqTemplate.js +++ b/msa/js-executor/queue/rabbitmqTemplate.js @@ -26,7 +26,7 @@ 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 queueProperties = config.get('rabbitmq.queue_properties'); const poolInterval = config.get('js.response_poll_interval'); const amqp = require('amqplib/callback_api'); diff --git a/msa/js-executor/queue/serviceBusTemplate.js b/msa/js-executor/queue/serviceBusTemplate.js index 034921afb7..20cf664940 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; 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.'); From e5c5aa705fa308397536c18c0779480d849ab368 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Thu, 30 Apr 2020 20:14:34 +0300 Subject: [PATCH 2/4] refactored --- docker/tb-js-executor.env | 2 +- msa/js-executor/queue/awsSqsTemplate.js | 12 +++++++----- 2 files changed, 8 insertions(+), 6 deletions(-) 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/queue/awsSqsTemplate.js b/msa/js-executor/queue/awsSqsTemplate.js index a5338c5e73..e0ffb55c87 100644 --- a/msa/js-executor/queue/awsSqsTemplate.js +++ b/msa/js-executor/queue/awsSqsTemplate.js @@ -74,11 +74,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(); From eedb38384536215555c97cb4fcff53ae7ee1bf72 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Fri, 1 May 2020 14:15:31 +0300 Subject: [PATCH 3/4] AWS improvements --- .../server/queue/sqs/TbAwsSqsConsumerTemplate.java | 6 +++++- .../server/queue/sqs/TbAwsSqsProducerTemplate.java | 5 ++++- .../server/queue/sqs/TbAwsSqsQueueAttributes.java | 1 - msa/js-executor/package.json | 1 + msa/js-executor/queue/awsSqsTemplate.js | 5 +++-- msa/js-executor/queue/pubSubTemplate.js | 8 +++++--- 6 files changed, 18 insertions(+), 8 deletions(-) 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 3e71388844..b66cad1504 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 @@ -127,6 +127,11 @@ public class TbAwsSqsConsumerTemplate implements TbQueueCo if (!subscribed) { List topicNames = partitions.stream().map(TopicPartitionInfo::getFullTopicName).collect(Collectors.toList()); queueUrls = topicNames.stream().map(this::getQueueUrl).collect(Collectors.toSet()); + + if (consumerExecutor != null) { + consumerExecutor.shutdown(); + } + consumerExecutor = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(queueUrls.size() * sqsSettings.getThreadsPerTopic() + 1)); subscribed = true; } @@ -172,7 +177,6 @@ public class TbAwsSqsConsumerTemplate implements TbQueueCo 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/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 e0ffb55c87..5f95de7d32 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'); @@ -29,7 +30,7 @@ const AWS = require('aws-sdk'); const queueProperties = config.get('aws_sqs.queue_properties'); const poolInterval = 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,7 @@ 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) { diff --git a/msa/js-executor/queue/pubSubTemplate.js b/msa/js-executor/queue/pubSubTemplate.js index 7d0b32ea34..cc5284022d 100644 --- a/msa/js-executor/queue/pubSubTemplate.js +++ b/msa/js-executor/queue/pubSubTemplate.js @@ -60,9 +60,11 @@ function PubSubProducer() { const topicList = await pubSubClient.getTopics(); if (topicList) { - topicList[0].forEach(topic => { - topics.push(getName(topic.name)); - }); + if (topicList) { + topicList[0].forEach(topic => { + topics.push(getName(topic.name)); + }); + } } const subscriptionList = await pubSubClient.getSubscriptions(); From 06c3caf082ee48eed660927ff453d530e965bfb6 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Sat, 2 May 2020 13:34:11 +0300 Subject: [PATCH 4/4] refactored --- .../server/queue/pubsub/TbPubSubAdmin.java | 32 ++++++++++++----- .../pubsub/TbPubSubConsumerTemplate.java | 5 +++ .../pubsub/TbPubSubProducerTemplate.java | 4 +-- msa/js-executor/queue/pubSubTemplate.js | 35 +++++++++++-------- pom.xml | 2 +- 5 files changed, 52 insertions(+), 26 deletions(-) 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 5dc795739a..495c895a91 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 @@ -134,6 +134,11 @@ public class TbPubSubConsumerTemplate implements TbQueueCo if (!subscribed) { subscriptionNames = partitions.stream().map(TopicPartitionInfo::getFullTopicName).collect(Collectors.toSet()); subscriptionNames.forEach(admin::createTopicIfNotExists); + + if (consumerExecutor != null) { + consumerExecutor.shutdown(); + } + consumerExecutor = Executors.newFixedThreadPool(subscriptionNames.size()); messagesPerTopic = pubSubSettings.getMaxMessages() / subscriptionNames.size(); subscribed = true; 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..7a073616fd 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 @@ -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/msa/js-executor/queue/pubSubTemplate.js b/msa/js-executor/queue/pubSubTemplate.js index cc5284022d..8cb264de64 100644 --- a/msa/js-executor/queue/pubSubTemplate.js +++ b/msa/js-executor/queue/pubSubTemplate.js @@ -60,11 +60,9 @@ function PubSubProducer() { const topicList = await pubSubClient.getTopics(); if (topicList) { - if (topicList) { - topicList[0].forEach(topic => { - topics.push(getName(topic.name)); - }); - } + topicList[0].forEach(topic => { + topics.push(getName(topic.name)); + }); } const subscriptionList = await pubSubClient.getSubscriptions(); @@ -100,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/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