diff --git a/src/kafka-pubsub.ts b/src/kafka-pubsub.ts index 60015f8..cad223d 100644 --- a/src/kafka-pubsub.ts +++ b/src/kafka-pubsub.ts @@ -113,7 +113,7 @@ export class KafkaPubSub implements PubSubEngine { {}, { topics: [topic] } ); - stream.consumer.on('data', (message) => { + stream.on('data', (message) => { let parsedMessage = JSON.parse(message.value.toString()) // Using channel abstraction