diff --git a/packages/kafka/lib/AbstractKafkaConsumer.ts b/packages/kafka/lib/AbstractKafkaConsumer.ts index 6f0557ce..2e44b2ca 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 not valid message(s)') 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 */ 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" },