diff --git a/msa/js-executor/api/jsInvokeMessageProcessor.js b/msa/js-executor/api/jsInvokeMessageProcessor.js index 867dd82b39..dec935c8ad 100644 --- a/msa/js-executor/api/jsInvokeMessageProcessor.js +++ b/msa/js-executor/api/jsInvokeMessageProcessor.js @@ -25,6 +25,7 @@ const config = require('config'), Utils = require('./utils'), JsExecutor = require('./jsExecutor'); +const statFrequency = Number(config.get('script.stat_print_frequency')); const scriptBodyTraceFrequency = Number(config.get('script.script_body_trace_frequency')); const useSandbox = config.get('script.use_sandbox') === 'true'; const maxActiveScripts = Number(config.get('script.max_active_scripts')); @@ -39,6 +40,7 @@ function JsInvokeMessageProcessor(producer) { this.scriptMap = new Map(); this.scriptIds = []; this.executedScriptsCounter = 0; + this.lastStatTime = performance.now(); } JsInvokeMessageProcessor.prototype.onJsInvokeMessage = function(message) { @@ -118,11 +120,16 @@ JsInvokeMessageProcessor.prototype.processInvokeRequest = function(requestId, re var scriptId = getScriptId(invokeRequest); logger.debug('[%s] Processing invoke request, scriptId: [%s]', requestId, scriptId); this.executedScriptsCounter++; - if ( this.executedScriptsCounter >= scriptBodyTraceFrequency ) { - this.executedScriptsCounter = 0; - if (logger.levels[logger.level] >= logger.levels['debug']) { - logger.debug('[%s] Executing script body: [%s]', scriptId, invokeRequest.scriptBody); - } + if (this.executedScriptsCounter % statFrequency == 0) { + var nowMs = performance.now(); + var msSinceLastStat = nowMs - this.lastStatTime; + var requestsPerSec = statFrequency / msSinceLastStat * 1000; //msSinceLastStat can't be zero in the real world + this.lastStatTime = nowMs; + logger.info('STAT[%s]: requests [%s], took [%s]ms, request/s [%s]', this.executedScriptsCounter, statFrequency, msSinceLastStat, requestsPerSec); + } + + if (this.executedScriptsCounter % scriptBodyTraceFrequency == 0) { + logger.info('[%s] Executing script body: [%s]', scriptId, invokeRequest.scriptBody); } this.getOrCompileScript(scriptId, invokeRequest.scriptBody).then( (script) => { diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index 5e2bfe8342..c5eb15e8b9 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -26,6 +26,8 @@ kafka: servers: "TB_KAFKA_SERVERS" replication_factor: "TB_QUEUE_KAFKA_REPLICATION_FACTOR" acks: "TB_KAFKA_ACKS" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge + batch_size: "${TB_KAFKA_BATCH_SIZE:128}" # for producer + linger_ms: "${TB_KAFKA_LINGER_MS:1}" # for producer requestTimeout: "TB_QUEUE_KAFKA_REQUEST_TIMEOUT_MS" compression: "TB_QUEUE_KAFKA_COMPRESSION" # gzip or uncompressed topic_properties: "TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES" @@ -70,6 +72,7 @@ logger: script: use_sandbox: "SCRIPT_USE_SANDBOX" + stat_print_frequency: "SCRIPT_STAT_PRINT_FREQUENCY" script_body_trace_frequency: "SCRIPT_BODY_TRACE_FREQUENCY" max_active_scripts: "MAX_ACTIVE_SCRIPTS" slow_query_log_ms: "SLOW_QUERY_LOG_MS" #1.123456 diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index cdf5b35421..5a2c47a675 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -26,6 +26,8 @@ kafka: servers: "localhost:9092" replication_factor: "1" acks: "1" # -1 = all; 0 = no acknowledgments; 1 = only waits for the leader to acknowledge + batch_size: "128" # for producer + linger_ms: "1" # for producer requestTimeout: "30000" # The default value in kafkajs is: 30000 compression: "gzip" # gzip or uncompressed topic_properties: "retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600;partitions:100;min.insync.replicas:1" @@ -59,7 +61,8 @@ logger: script: use_sandbox: "true" - script_body_trace_frequency: "1000" + script_body_trace_frequency: "10000" + stat_print_frequency: "10000" max_active_scripts: "1000" - slow_query_log_ms: "1.000000" #millis + slow_query_log_ms: "5.000000" #millis slow_query_log_body: "false" diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index a061215b28..79d87c17b4 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -23,12 +23,11 @@ const replicationFactor = Number(config.get('kafka.replication_factor')); const topicProperties = config.get('kafka.topic_properties'); const kafkaClientId = config.get('kafka.client_id'); const acks = Number(config.get('kafka.acks')); +const maxBatchSize = Number(config.get('kafka.batch_size')); +const linger = Number(config.get('kafka.linger_ms')); const requestTimeout = Number(config.get('kafka.requestTimeout')); const compressionType = (config.get('kafka.requestTimeout') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; -const linger = 5; //milliseconds //TODO move to the config -const maxBatchSize = 10; //max messages in batch //TODO move to the config - let kafkaClient; let kafkaAdmin; let consumer; @@ -76,7 +75,9 @@ function sendLoopWithLinger() { function sendMessagesAsBatch() { if (batchMessages.length > 0) { - logger.info('sendMessagesAsBatch, lenght: [%s]', batchMessages.length ); + if (batchMessages.length > 1) { + logger.info('sendMessagesAsBatch, length: [%s]', batchMessages.length); + } const messagesToSend = batchMessages; batchMessages = []; producer.sendBatch({ @@ -85,17 +86,15 @@ function sendMessagesAsBatch() { compression: compressionType }).then( () => { - logger.info('Response sent to kafka, length: [%s]', messagesToSend.length ); + 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 ); + logger.error('Failed to send kafka, length: [%s], pending to reprocess msgs', messagesToSend.length); batchMessages = messagesToSend.concat(batchMessages); logger.error(err.stack); } ); - } else { - //logger.debug("nothing to send"); } } @@ -163,6 +162,8 @@ function sendMessagesAsBatch() { consumer = kafkaClient.consumer({groupId: 'js-executor-group'}); producer = kafkaClient.producer(); +/* + //producer event instrumentation to debug const { CONNECT } = producer.events; const removeListenerC = producer.on(CONNECT, e => logger.info(`producer CONNECT`)); const { DISCONNECT } = producer.events; @@ -173,18 +174,22 @@ function sendMessagesAsBatch() { const removeListenerRT = producer.on(REQUEST_TIMEOUT, e => logger.info(`producer REQUEST_TIMEOUT ${e.payload.broker}`)); const { REQUEST_QUEUE_SIZE } = producer.events; const removeListenerRQS = producer.on(REQUEST_QUEUE_SIZE, e => logger.info(`producer REQUEST_QUEUE_SIZE ${e.payload.broker} size ${e.queueSize}`)); +*/ -// const removeListeners = {} -// const { FETCH_START } = consumer.events; -// removeListeners[FETCH_START] = consumer.on(FETCH_START, e => logger.info(`consumer FETCH_START`)); -// const { FETCH } = consumer.events; -// removeListeners[FETCH] = consumer.on(FETCH, e => logger.info(`consumer FETCH numberOfBatches ${e.payload.numberOfBatches} duration ${e.payload.duration}`)); -// const { START_BATCH_PROCESS } = consumer.events; -// removeListeners[START_BATCH_PROCESS] = consumer.on(START_BATCH_PROCESS, e => logger.info(`consumer START_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); -// const { END_BATCH_PROCESS } = consumer.events; -// removeListeners[END_BATCH_PROCESS] = consumer.on(END_BATCH_PROCESS, e => logger.info(`consumer END_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); -// const { COMMIT_OFFSETS } = consumer.events; -// removeListeners[COMMIT_OFFSETS] = consumer.on(COMMIT_OFFSETS, e => logger.info(`consumer COMMIT_OFFSETS topics ${e.payload.topics}`)); +/* + //consumer event instrumentation to debug + const removeListeners = {} + const { FETCH_START } = consumer.events; + removeListeners[FETCH_START] = consumer.on(FETCH_START, e => logger.info(`consumer FETCH_START`)); + const { FETCH } = consumer.events; + removeListeners[FETCH] = consumer.on(FETCH, e => logger.info(`consumer FETCH numberOfBatches ${e.payload.numberOfBatches} duration ${e.payload.duration}`)); + const { START_BATCH_PROCESS } = consumer.events; + removeListeners[START_BATCH_PROCESS] = consumer.on(START_BATCH_PROCESS, e => logger.info(`consumer START_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); + const { END_BATCH_PROCESS } = consumer.events; + removeListeners[END_BATCH_PROCESS] = consumer.on(END_BATCH_PROCESS, e => logger.info(`consumer END_BATCH_PROCESS topic ${e.payload.topic} batchSize ${e.payload.batchSize}`)); + const { COMMIT_OFFSETS } = consumer.events; + removeListeners[COMMIT_OFFSETS] = consumer.on(COMMIT_OFFSETS, e => logger.info(`consumer COMMIT_OFFSETS topics ${e.payload.topics}`)); +*/ const messageProcessor = new JsInvokeMessageProcessor(new KafkaProducer()); await consumer.connect();