Browse Source

js-executor: added parameters for producer TB_KAFKA_BATCH_SIZE and TB_KAFKA_LINGER_MS; added print stats frequency SCRIPT_STAT_PRINT_FREQUENCY

pull/4667/head
Sergey Matvienko 5 years ago
parent
commit
93bea70205
  1. 17
      msa/js-executor/api/jsInvokeMessageProcessor.js
  2. 3
      msa/js-executor/config/custom-environment-variables.yml
  3. 7
      msa/js-executor/config/default.yml
  4. 43
      msa/js-executor/queue/kafkaTemplate.js

17
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) => {

3
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

7
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"

43
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();

Loading…
Cancel
Save