From 947e75a56adf1edbed79d86ca0eae5fececa1751 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Thu, 7 May 2026 13:33:46 +0300 Subject: [PATCH] Code review changes --- msa/js-executor/package.json | 5 +++-- msa/js-executor/queue/kafkaTemplate.ts | 19 +++++++++++++------ 2 files changed, 16 insertions(+), 8 deletions(-) diff --git a/msa/js-executor/package.json b/msa/js-executor/package.json index dfa50f6765..7cbd99ea48 100644 --- a/msa/js-executor/package.json +++ b/msa/js-executor/package.json @@ -13,11 +13,11 @@ "build": "tsc" }, "dependencies": { + "@2l/kafkajs-lz4": "^1.3.2", "config": "^4.1.1", "express": "^5.1.0", "js-yaml": "^4.1.1", "kafkajs": "^2.2.4", - "@2l/kafkajs-lz4": "^1.3.2", "long": "^5.3.2", "uuid-parse": "^1.1.0", "winston": "^3.17.0", @@ -47,7 +47,8 @@ }, "pkg": { "assets": [ - "node_modules/config/**/*.*" + "node_modules/config/**/*.*", + "node_modules/@antoniomuso/lz4-napi-*/**/*.node" ] } } diff --git a/msa/js-executor/queue/kafkaTemplate.ts b/msa/js-executor/queue/kafkaTemplate.ts index 5a340b2cfd..1114190e64 100644 --- a/msa/js-executor/queue/kafkaTemplate.ts +++ b/msa/js-executor/queue/kafkaTemplate.ts @@ -31,14 +31,11 @@ import { Producer, TopicMessages } from 'kafkajs'; -import LZ4Codec from '@2l/kafkajs-lz4'; import { isNotEmptyStr } from '../api/utils'; import { KeyObject } from 'tls'; import process, { exit, kill } from 'process'; -CompressionCodecs[CompressionTypes.LZ4] = new LZ4Codec().codec; - export class KafkaTemplate implements IQueue { private logger = _logger(`kafkaTemplate`); @@ -50,15 +47,25 @@ export class KafkaTemplate implements IQueue { private linger = Number(config.get('kafka.linger_ms')); private requestTimeout = Number(config.get('kafka.requestTimeout')); private connectionTimeout = Number(config.get('kafka.connectionTimeout')); - private compressionType = KafkaTemplate.resolveCompressionType(config.get('kafka.compression')); + private compressionType = this.resolveCompressionType(config.get('kafka.compression')); - private static resolveCompressionType(compression: string): CompressionTypes { + private resolveCompressionType(compression: string): CompressionTypes { switch (compression) { case 'gzip': return CompressionTypes.GZIP; - case 'lz4': + case 'lz4': { + // Load the LZ4 codec lazily so users who don't enable LZ4 don't take a hard + // dependency on the lz4-napi native binary (e.g. inside pkg-built executables). + const LZ4Codec = require('@2l/kafkajs-lz4').default; + CompressionCodecs[CompressionTypes.LZ4] = new LZ4Codec().codec; return CompressionTypes.LZ4; + } + case 'none': + return CompressionTypes.None; default: + if (isNotEmptyStr(compression)) { + this.logger.warn('Unknown kafka.compression value "%s"; falling back to no compression. Supported values: gzip, lz4, none.', compression); + } return CompressionTypes.None; } }