From 1f2715093c5ba4410a7c52476f433ed11d2e0683 Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Wed, 3 Dec 2025 17:57:57 +0100 Subject: [PATCH 1/6] Improving kafka consumer logging --- packages/kafka/lib/AbstractKafkaConsumer.ts | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 6f0557ce..24a47824 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -262,6 +262,7 @@ export abstract class AbstractKafkaConsumer< messageOrBatch: MessageOrBatch>, ): Promise { const messageProcessingStartTimestamp = Date.now() + this.logger.debug({ origin: this.constructor.name, topic }, 'Consuming message(s)') const handlerConfig = this.resolveHandler(topic) @@ -275,12 +276,17 @@ export abstract class AbstractKafkaConsumer< ) if (!validMessages.length) { + this.logger.debug({ origin: this.constructor.name, topic }, 'Received message(s) is empty') return this.commit(messageOrBatch) + } else { + this.logger.debug( + { origin: this.constructor.name, topic, validMessagesCount: validMessages.length }, + 'Received valid message(s) to process', + ) } // biome-ignore lint/style/noNonNullAssertion: we check validMessages length above const firstMessage = validMessages[0]! - const requestContext = this.getRequestContext(firstMessage) /* v8 ignore next */ From 830914a37d601c7539ad3c996c7fd5044098efba Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Wed, 3 Dec 2025 17:58:47 +0100 Subject: [PATCH 2/6] Release prepare --- packages/kafka/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/kafka/package.json b/packages/kafka/package.json index 86fdec15..27196de5 100644 --- a/packages/kafka/package.json +++ b/packages/kafka/package.json @@ -1,6 +1,6 @@ { "name": "@message-queue-toolkit/kafka", - "version": "0.7.8", + "version": "0.7.9", "engines": { "node": ">= 22.14.0" }, From 3f1758cd43b50f8ccf71162310e42da52769b6ef Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Wed, 3 Dec 2025 18:01:40 +0100 Subject: [PATCH 3/6] Fix log --- packages/kafka/lib/AbstractKafkaConsumer.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 24a47824..a697f3c5 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -276,7 +276,7 @@ export abstract class AbstractKafkaConsumer< ) if (!validMessages.length) { - this.logger.debug({ origin: this.constructor.name, topic }, 'Received message(s) is empty') + this.logger.debug({ origin: this.constructor.name, topic }, 'Received not valid message(s)') return this.commit(messageOrBatch) } else { this.logger.debug( From e55dde285079aeb44af6d0cdf5b8b6635a100149 Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Wed, 3 Dec 2025 18:11:45 +0100 Subject: [PATCH 4/6] Lint fix --- packages/kafka/lib/AbstractKafkaConsumer.ts | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index a697f3c5..45a41cd8 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -262,7 +262,10 @@ export abstract class AbstractKafkaConsumer< messageOrBatch: MessageOrBatch>, ): Promise { const messageProcessingStartTimestamp = Date.now() - this.logger.debug({ origin: this.constructor.name, topic }, 'Consuming message(s)') + this.logger.debug( + { origin: this.constructor.name, topic, count: messageOrBatch }, + 'Consuming message(s)', + ) const handlerConfig = this.resolveHandler(topic) @@ -279,10 +282,10 @@ export abstract class AbstractKafkaConsumer< this.logger.debug({ origin: this.constructor.name, topic }, 'Received not valid message(s)') return this.commit(messageOrBatch) } else { - this.logger.debug( + this.logger.debug( { origin: this.constructor.name, topic, validMessagesCount: validMessages.length }, 'Received valid message(s) to process', - ) + ) } // biome-ignore lint/style/noNonNullAssertion: we check validMessages length above From de010f629d1c8ac40b734b97048d05dc8a065d70 Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Wed, 3 Dec 2025 19:45:27 +0100 Subject: [PATCH 5/6] AI suggestion --- packages/kafka/lib/AbstractKafkaConsumer.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 45a41cd8..556a88ca 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -263,7 +263,7 @@ export abstract class AbstractKafkaConsumer< ): Promise { const messageProcessingStartTimestamp = Date.now() this.logger.debug( - { origin: this.constructor.name, topic, count: messageOrBatch }, + { origin: this.constructor.name, topic }, 'Consuming message(s)', ) From 81430e3e6b692adc5efe9ee75f28796896bd1182 Mon Sep 17 00:00:00 2001 From: CarlosGamero Date: Wed, 3 Dec 2025 19:50:46 +0100 Subject: [PATCH 6/6] lint --- packages/kafka/lib/AbstractKafkaConsumer.ts | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 556a88ca..2e44b2ca 100644 --- a/packages/kafka/lib/AbstractKafkaConsumer.ts +++ b/packages/kafka/lib/AbstractKafkaConsumer.ts @@ -262,10 +262,7 @@ export abstract class AbstractKafkaConsumer< messageOrBatch: MessageOrBatch>, ): Promise { const messageProcessingStartTimestamp = Date.now() - this.logger.debug( - { origin: this.constructor.name, topic }, - 'Consuming message(s)', - ) + this.logger.debug({ origin: this.constructor.name, topic }, 'Consuming message(s)') const handlerConfig = this.resolveHandler(topic)