|
|
@ -42,7 +42,7 @@ export class KafkaTemplate implements IQueue { |
|
|
private maxBatchSize = Number(config.get('kafka.batch_size')); |
|
|
private maxBatchSize = Number(config.get('kafka.batch_size')); |
|
|
private linger = Number(config.get('kafka.linger_ms')); |
|
|
private linger = Number(config.get('kafka.linger_ms')); |
|
|
private requestTimeout = Number(config.get('kafka.requestTimeout')); |
|
|
private requestTimeout = Number(config.get('kafka.requestTimeout')); |
|
|
private connectionTimeout = Number(config.get('kafka.connection_timeout_ms')); |
|
|
private connectionTimeout = Number(config.get('kafka.connectionTimeout')); |
|
|
private compressionType = (config.get('kafka.compression') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; |
|
|
private compressionType = (config.get('kafka.compression') === "gzip") ? CompressionTypes.GZIP : CompressionTypes.None; |
|
|
private partitionsConsumedConcurrently = Number(config.get('kafka.partitions_consumed_concurrently')); |
|
|
private partitionsConsumedConcurrently = Number(config.get('kafka.partitions_consumed_concurrently')); |
|
|
|
|
|
|
|
|
|