KafkaJS消费Kafka消息无法读取 报ReferenceError问题解决
使用KafkaJS消费AWS MSK消息时无法进入eachMessage回调正常读取消息,运行时抛出如下错误:
{"level":"ERROR","timestamp":"2022-06-29T04:54:57.414Z","logger":"kafkajs","message":"[Runner] Error when calling eachMessage","topic":"aws.identity.users.0","partition":7,"offset":"52","stack":"ReferenceError: consumedMessages is not defined\n at Runner.eachMessage (C:\biport\testkafka\index.js:24:3)\n at Runner.processEachMessage (C:\biport\testkafka\node_modules\kafkajs\src\consumer\runner.js:191:20)\n at onBatch (C:\biport\testkafka\node_modules\kafkajs\src\consumer\runner.js:393:20)\n at C:\biport\testkafka\node_modules\kafkajs\src\consumer\runner.js:409:15\n at retry (C:\biport\testkafka\node_modules\kafkajs\src\retry\index.js:43:5)\n at C:\biport\testkafka\node_modules\kafkajs\src\retry\index.js:75:5)\n at new Promise ()\n at Runner.retrier (C:\biport\testkafka\node_modules\kafkajs\src\retry\index.js:72:10)\n at Runner.handleBatch (C:\biport\testkafka\node_modules\kafkajs\src\consumer\runner.js:407:17)\n at handler (C:\biport\testkafka\node_modules\kafkajs\src\consumer\runner.js:58:30)","error":{}}
- 变量访问异常:错误核心是
ReferenceError: consumedMessages is not defined,eachMessage消费回调执行时无法访问外部定义的consumedMessages数组。一方面代码里两处裸写的无意义模板字符串(`not able to get inside this`、`how to make this code work`)属于无效代码,干扰了JS作用域解析;另一方面如果实际运行代码和贴出的代码不一致,变量定义位置错误也会触发该问题。 - 消费者启动逻辑位置错误:消费者启动的
run()方法被写在根路由/的处理函数中,服务启动时不会自动连接Kafka,只有收到对应GET请求才会初始化消费者,且多次触发请求会重复执行连接、订阅、消费启动逻辑,造成消费者实例状态混乱、连接泄漏,进一步导致回调执行异常。 - 配置存在隐患:连接的是AWS MSK集群,如果集群开启了IAM认证,仅配置
ssl: true缺少SASL认证参数会直接连不上集群;同时设置了autoCommit: false关闭自动提交偏移量,但消费逻辑里没有手动提交偏移量,消费者重启后会重复拉取历史消息;另外Kafka消息的key、value默认是Buffer二进制类型,直接打印看不到明文内容,很容易误以为没有读到消息。 - 基础代码缺失:代码中
app.listen用到的port变量没有定义,会直接导致服务启动失败。
- 清理所有无效的裸写模板字符串,确认
consumedMessages变量定义在消费回调可访问的上层作用域,解决引用报错。 - 把消费者启动逻辑从路由处理函数中移出,改为服务启动后自动执行,增加启动状态标记避免重复连接。
- 补全缺失的配置:如果使用AWS MSK IAM认证,补充对应SASL认证参数;关闭自动提交的场景下,消息处理完成后手动提交偏移量;把Buffer类型的消息key、value转为字符串,方便查看消费到的内容。
- 定义缺失的
port变量,补全必要的模块引入,保证服务可以正常启动。
const { Kafka, logLevel } = require('kafkajs') const express = require('express') const app = express() const port = 3000 // 补全缺失的端口定义 // 若使用AWS MSK IAM认证,需要引入对应认证依赖并补充sasl配置,非IAM认证可忽略 // const { MSKAuthMechanism } = require('@aws-sdk/kafka-sasl-iam') const kafka = new Kafka({ logLevel: logLevel.INFO, ssl: true, // IAM认证时打开下面的配置 // sasl: { // mechanism: MSKAuthMechanism, // authenticationTimeout: 10000, // }, brokers: [ 'b-1.dev-datafabric.nmz564.c20.kafka.us-east-1.amazonaws.com:9094', 'b-2.dev-datafabric.nmz564.c20.kafka.us-east-1.amazonaws.com:9094', 'b-3.dev-datafabric.nmz564.c20.kafka.us-east-1.amazonaws.com:9094' ], clientId: 'local-client' }) const topic = 'aws.identity.users.0' const consumer = kafka.consumer({ groupId: 'test-group' }) const consumedMessages = [] let consumerStarted = false // 启动标记,防止重复连接 const runConsumer = async () => { if (consumerStarted) return consumerStarted = true await consumer.connect() await consumer.subscribe({ topic, fromBeginning: true }) await consumer.run({ autoCommit: false, eachMessage: async ({ topic, partition, message }) => { consumedMessages.push(message); // 将Buffer类型的消息内容转为字符串,方便查看 const msgKey = message.key ? message.key.toString() : null const msgValue = message.value ? message.value.toString() : null const prefix = `${topic}[${partition} | ${message.offset}] / ${message.timestamp}` console.log(`- ${prefix} ${msgKey}#${msgValue}`) // 处理完消息后手动提交偏移量,避免重复消费 await consumer.commitOffsets([{ topic, partition, offset: (Number(message.offset) + 1).toString() }]) }, }) } app.get('/', (req, res) => { res.send('Hello World!') }); // 服务启动时自动初始化消费者 app.listen(port, async () => { console.log(`Example app listening on port ${port}!`) try { await runConsumer() console.log('Kafka consumer started successfully') } catch (e) { console.error(`[example/consumer] ${e.message}`, e) } });
内容的提问来源于stack exchange,提问作者Nandish Kumar

