You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 12:57:11