From 9d82899f049b44df2c9a2d51f584d840427c00be Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Sun, 27 Sep 2020 12:41:49 +0300 Subject: [PATCH] raised kafkajs version and improvements kafkaTemplate.js --- msa/js-executor/package.json | 2 +- msa/js-executor/queue/kafkaTemplate.js | 19 ++++++------------- 2 files changed, 7 insertions(+), 14 deletions(-) diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index 967f55b24f..405d582f04 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -14,7 +14,7 @@ "dependencies": { "config": "^3.2.2", "js-yaml": "^3.12.0", - "kafkajs": "^1.12.0", + "kafkajs": "^1.14.0", "@google-cloud/pubsub": "^1.7.1", "aws-sdk": "^2.663.0", "amqplib": "^0.5.5", diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 2672a8df6a..33dd0d8c20 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -27,20 +27,10 @@ let kafkaAdmin; let consumer; let producer; -const topics = []; const configEntries = []; function KafkaProducer() { this.send = async (responseTopic, scriptId, rawResponse, headers) => { - - if (!topics.includes(responseTopic)) { - let createResponseTopicResult = await createTopic(responseTopic, 1); - topics.push(responseTopic); - if (createResponseTopicResult) { - logger.info('Created new topic: %s', requestTopic); - } - } - return producer.send( { topic: responseTopic, @@ -99,10 +89,13 @@ function KafkaProducer() { } } - let createRequestTopicResult = await createTopic(requestTopic, partitions); + let topics = await kafkaAdmin.listTopics(); - if (createRequestTopicResult) { - logger.info('Created new topic: %s', requestTopic); + if (!topics.includes(requestTopic)) { + let createRequestTopicResult = await createTopic(requestTopic, partitions); + if (createRequestTopicResult) { + logger.info('Created new topic: %s', requestTopic); + } } consumer = kafkaClient.consumer({groupId: 'js-executor-group'});