From 2d6ee5ff2cde7dedabff9b65bdcd233f04b710b2 Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Thu, 11 Dec 2025 13:50:33 +0100 Subject: [PATCH 1/3] Waiting for a batch to be processed before accepting more messages --- packages/kafka/lib/AbstractKafkaConsumer.ts | 16 +++++----------- .../kafka/lib/utils/KafkaMessageBatchStream.ts | 2 +- 2 files changed, 6 insertions(+), 12 deletions(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 2e44b2ca..9e1bc225 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -206,24 +206,18 @@ export abstract class AbstractKafkaConsumer< } if (this.options.batchProcessingEnabled && this.messageBatchStream) { - this.messageBatchStream.on('data', async (messageBatch) => - this.consume(messageBatch.topic, messageBatch.messages), - ) - this.messageBatchStream.on('error', (error) => this.handlerError(error)) + this.handleSyncStream(this.messageBatchStream).catch((error) => this.handlerError(error)) } else { - // biome-ignore lint/style/noNonNullAssertion: consumerStream is always created - const stream = this.consumerStream! - - // we are not waiting for the stream to complete - // because init() must return promised void - this.handleSyncStream(stream).catch((error) => this.handlerError(error)) + this.handleSyncStream(this.consumerStream).catch((error) => this.handlerError(error)) } this.consumerStream.on('error', (error) => this.handlerError(error)) } private async handleSyncStream( - stream: MessagesStream, + stream: + | MessagesStream + | KafkaMessageBatchStream>>, ): Promise { for await (const message of stream) { await this.consume( diff --git a/packages/kafka/lib/utils/KafkaMessageBatchStream.ts b/packages/kafka/lib/utils/KafkaMessageBatchStream.ts index 54178c0a..5ae29df8 100644 --- a/packages/kafka/lib/utils/KafkaMessageBatchStream.ts +++ b/packages/kafka/lib/utils/KafkaMessageBatchStream.ts @@ -14,7 +14,7 @@ export type MessageBatch = { topic: string; partition: number; message export interface KafkaMessageBatchStream extends Duplex { - // biome-ignore lint/suspicious/noExplicitAny: compatible with Duplex definition + // biome-ignore lint/suspicious/noExplicitAny: compatible with Duplex definition on(event: string | symbol, listener: (...args: any[]) => void): this on(event: 'data', listener: (chunk: MessageBatch) => void): this From 006dc0904e18db9369197df29c2196637d5ce135 Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Thu, 11 Dec 2025 13:52:30 +0100 Subject: [PATCH 2/3] Release prepare + dep update --- packages/kafka/package.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/kafka/package.json b/packages/kafka/package.json index 753fe344..cf7bd405 100644 --- a/packages/kafka/package.json +++ b/packages/kafka/package.json @@ -1,6 +1,6 @@ { "name": "@message-queue-toolkit/kafka", - "version": "0.8.0", + "version": "0.8.1", "engines": { "node": ">= 22.14.0" }, @@ -53,7 +53,7 @@ "dependencies": { "@lokalise/node-core": "^14.2.0", "@lokalise/universal-ts-utils": "^4.5.1", - "@platformatic/kafka": "^1.21.0" + "@platformatic/kafka": "^1.22.0" }, "peerDependencies": { "@message-queue-toolkit/core": ">=23.0.0", From 7b2b252d39bd13a99a48a6613d1ccba41c2cc7cf Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Thu, 11 Dec 2025 14:07:43 +0100 Subject: [PATCH 3/3] Fixing type issue --- packages/kafka/lib/AbstractKafkaConsumer.ts | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 9e1bc225..0e5764af 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -206,7 +206,7 @@ export abstract class AbstractKafkaConsumer< } if (this.options.batchProcessingEnabled && this.messageBatchStream) { - this.handleSyncStream(this.messageBatchStream).catch((error) => this.handlerError(error)) + this.handleSyncStreamBatch(this.messageBatchStream).catch((error) => this.handlerError(error)) } else { this.handleSyncStream(this.consumerStream).catch((error) => this.handlerError(error)) } @@ -215,9 +215,7 @@ export abstract class AbstractKafkaConsumer< } private async handleSyncStream( - stream: - | MessagesStream - | KafkaMessageBatchStream>>, + stream: MessagesStream, ): Promise { for await (const message of stream) { await this.consume( @@ -226,6 +224,16 @@ export abstract class AbstractKafkaConsumer< ) } } + private async handleSyncStreamBatch( + stream: KafkaMessageBatchStream>>, + ): Promise { + for await (const messageBatch of stream) { + await this.consume( + messageBatch.topic, + messageBatch.messages as DeserializedMessage>, + ) + } + } async close(): Promise { if (!this.consumerStream && !this.messageBatchStream) {