From a4e28ad94533610c4de0f855bad75dcd4460ac74 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 1 Jun 2021 16:00:32 +0300 Subject: [PATCH] js-executor fixed promises for each message for Kafka batches --- .../api/jsInvokeMessageProcessor.js | 24 +++++++------- msa/js-executor/queue/kafkaTemplate.js | 33 ++++++++++++------- 2 files changed, 33 insertions(+), 24 deletions(-) diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index dec935c8ad..1aabd216d2 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -180,19 +180,17 @@ JsInvokeMessageProcessor.prototype.sendResponse = function (requestId, responseT var remoteResponse = createRemoteResponse(requestId, compileResponse, invokeResponse, releaseResponse); var rawResponse = Buffer.from(JSON.stringify(remoteResponse), 'utf8'); logger.debug('[%s] Sending response to queue, scriptId: [%s]', requestId, scriptId); - this.producer.send(responseTopic, scriptId, rawResponse, headers); -//TODO put error msg for other queues implementation except Kafka -// .then( -// () => { -// logger.info('[%s] Response sent to queue, took [%s]ms, scriptId: [%s]', requestId, (performance.now() - tStartSending), scriptId); -// }, -// (err) => { -// if (err) { -// logger.error('[%s] Failed to send response to queue: %s', requestId, err.message); -// logger.error(err.stack); -// } -// } -// ); + this.producer.send(responseTopic, scriptId, rawResponse, headers).then( + () => { + logger.debug('[%s] Response sent to queue, took [%s]ms, scriptId: [%s]', requestId, (performance.now() - tStartSending), scriptId); + }, + (err) => { + if (err) { + logger.error('[%s] Failed to send response to queue: %s', requestId, err.message); + logger.error(err.stack); + } + } + ); } JsInvokeMessageProcessor.prototype.getOrCompileScript = function(scriptId, scriptBody) { diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index f467d0813b..1fac350976 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -36,11 +36,12 @@ let producer; const configEntries = []; let batchMessages = []; +let batchResolvers = []; let sendLoopInstance; function KafkaProducer() { this.send = (responseTopic, scriptId, rawResponse, headers) => { - logger.debug('Pending queue response, scriptId: [%s]', scriptId); + logger.debug('Pending queue response, scriptId: [%s]', scriptId); const message = { topic: responseTopic, messages: [{ @@ -50,17 +51,22 @@ function KafkaProducer() { }] }; - pushMessageToSendLater(message); - return {}; + return pushMessageToSendLater(message); } } function pushMessageToSendLater(message) { + let resolver; + const promise = new Promise((resolve, reject) => { + resolver = resolve; + }); batchMessages.push(message); + batchResolvers.push(resolver); if (batchMessages.length >= maxBatchSize) { sendMessagesAsBatch(); sendLoopWithLinger(); //reset loop function and reschedule new linger } + return promise; } function sendLoopWithLinger() { @@ -77,22 +83,27 @@ function sendMessagesAsBatch() { if (batchMessages.length > 0) { logger.debug('sendMessagesAsBatch, length: [%s]', batchMessages.length); const messagesToSend = batchMessages; + const resolvers = batchResolvers; batchMessages = []; + batchResolvers = []; producer.sendBatch({ topicMessages: messagesToSend, acks: acks, compression: compressionType }).then( - () => { - logger.debug('Response sent to kafka, length: [%s]', messagesToSend.length); - }, - (err) => { - logger.error('Failed to send kafka, length: [%s], pending to reprocess msgs', messagesToSend.length); - batchMessages = messagesToSend.concat(batchMessages); - logger.error(err.stack); + () => { + logger.debug('Response batch sent to kafka, length: [%s]', messagesToSend.length); + for (let i = 0; i < promisesToSend.length; i++) { + resolvers[i](); } + }, + (err) => { + logger.error('Failed batch send to kafka, length: [%s], pending to reprocess msgs', messagesToSend.length); + logger.error(err.stack); + batchMessages = messagesToSend.concat(batchMessages); + batchResolvers = resolvers.concat(batchResolvers); //promises will never be rejected. Will retry forever + } ); - } }