AWS MSK(IAM认证)与Node.js应用Kafka连接报错求助
AWS MSK(IAM认证)Node.js消费者连接错误解决
问题概述
我有一个启用IAM认证的AWS MSK集群,使用Node.js的KafkaJS库连接时,消费者和Admin客户端出现KafkaJSConnectionClosedError: Closed connection错误。应用部署在MSK集群的同一VPC内,为测试暂时注释了生产者和消费者代码中的SASL配置。
生产者代码
const { Kafka, logLevel } = require('kafkajs'); const kafka = new Kafka({ clientId: 'api_gateway_client', brokers: ["b-1.ashique.exxxx.xxxxafka.ap-southexxxx.xxzonaws.com:9098"], ssl: true, // sasl: { // mechanism: 'aws', // authorizationIdentity: '882xxxx397', // UserId or RoleId // accessKeyId: 'AKIAxxxxxD37NEJ', // secretAccessKey: 'q7SI3JxxxxxxxMhY0ciAou6TDXWdyR6h', //sessionToken: 'WHArYt8i5vfxxxxxxv/Eww6eL9tgQMJp6QFNEXAMPLETOKEN' // Optional // }, }) const producer = kafka.producer(); producer.connect(); export async function sendKafkaMessage(topic: string, message: string): Promise<any> { const result = await producer.send({ topic, messages: [ { value: message } ] }); return result; } producer.on('producer.connect',() => { console.log('Kafka producer connected ................'); const admin = kafka.admin(); });
消费者代码
//import { Kafka, Consumer, EachMessagePayload } from 'kafkajs' const { Kafka, logLevel } = require('kafkajs'); //code for plain mechanism const kafka = new Kafka({ clientId: 'api_gateway_client', brokers: ["b-1.ashique.et7xxxxx.xxx.xxxonaws.com:9098"], ssl: true, // logLevel: logLevel.DEBUG, // sasl: { // mechanism: 'plain', // username: "AKIxxxxxD37NEJ", // password: "q7SI3JyxxxxxxciAou6TDXWdyR6h" // } // }); }); const admin = kafka.admin(); async function listTopics() { await admin.connect(); const topics = await admin.listTopics(); console.log('Topics:', topics); } type Ttopic = { Name: string, FromBeginning: boolean, } const topics: Ttopic[] = [ { Name: 'ashique',FromBeginning: false } // first create a topic in kafka and specify the topic here ] // Define the consumer options const consumerOptions = { groupId: 'api_gateway_group', maxWaitTimeInMs: 5000, } // Define the function that will handle incoming messages async function handleMessage({ topic, partition, message }: any) { console.log(`Received message on topic ${topic}, partition ${partition}: ${message.value.toString()}`) } // Create the KafkaJS consumer instance const consumer: any = kafka.consumer(consumerOptions) async function startConsumer() { // Connect to the Kafka brokers await consumer.connect() // Subscribe to the topics for (const topic of topics) { console.log("topicName ==",topic.Name); await consumer.subscribe({ topic: topic.Name, fromBeginning: topic.FromBeginning }); } await consumer.run({ eachMessage: handleMessage, }); } listTopics().catch(error => { console.error('Error listing topics:', error); }) // Start the consumer startConsumer().catch(error => { console.log('Error starting consumer::::', error); });
错误日志
{"level":"WARN","timestamp":"2023-08-26T06:00:45.361Z","logger":"kafkajs","message":"KafkaJS v2.0.0 switched default partitioner. To retain the same partitioning behavior as in previous versions, create the producer with the option \"createPartitioner: Partitioners.LegacyPartitioner\". See the migration guide at https://kafka.js.org/docs/migration-guide-v2.0.0#producer-new-default-partitioner for details. Silence this warning by setting the environment variable \"KAFKAJS_NO_PARTITIONER_WARNING=1\""} port running on - 5000 Kafka producer connected ................ topicName == ashique Error listing topics: KafkaJSConnectionClosedError: Closed connection at TLSSocket.onEnd (/var/www/html/KafkaTestCode/node_modules/kafkajs/src/network/connection.js:197:13) at TLSSocket.emit (node:events:525:35) at TLSSocket.emit (node:domain:489:12) at endReadableNT (node:internal/streams/readable:1359:12) at processTicksAndRejections (node:internal/process/task_queues:82:21) { retriable: true, helpUrl: undefined, broker: 'b-1.ashiquxxxxx.xxxx.xxxxxxazonaws.com:9098', code: undefined, host: 'b-1.ashiquxxxxx.xxxx.xxxxxxazonaws.com', port: 9098, [cause]: undefined } Error starting consumer:::: KafkaJSConnectionClosedError: Closed connection at TLSSocket.onEnd (/var/www/html/KafkaTestCode/node_modules/kafkajs/src/network/connection.js:197:13) at TLSSocket.emit (node:events:525:35) at TLSSocket.emit (node:domain:489:12) at endReadableNT (node:internal/streams/readable:1359:12) at processTicksAndRejections (node:internal/process/task_queues:82:21) { retriable: true, helpUrl: undefined, broker: 'b-1.ashiquxxxxx.xxxx.xxxxxxazonaws.com:9098', code: undefined, host: 'b-1.ashiquxxxxx.xxxx.xxxxxxazonaws.com', port: 9098, [cause]: undefined }
解决方案
- 必须启用IAM认证的SASL配置:MSK的9098端口是IAM认证专属端口,不允许无认证连接。注释SASL配置会导致集群主动关闭连接,这是错误的核心原因。
- 统一并修复Broker地址:生产者和消费者代码中的Broker地址不一致,需修改为完全相同的MSK集群Broker地址。
- 正确配置IAM SASL参数:
- 生产者和消费者的Kafka实例都要配置
sasl: { mechanism: 'aws', ... },不能仅配置一端 - 优先使用IAM角色而非硬编码密钥:如果应用部署在EC2/EKS等AWS服务上,给实例/Pod附加具备MSK访问权限的IAM角色,KafkaJS会自动从环境获取凭证,无需硬编码
accessKeyId和secretAccessKey - 确保IAM实体(用户/角色)拥有必要权限:通过IAM策略配置
kafka:DescribeCluster、kafka:ListTopics、kafka:Subscribe、kafka:Consume等权限
- 生产者和消费者的Kafka实例都要配置
- 修复Admin客户端使用:消费者代码中
listTopics函数连接Admin后未断开,建议在函数末尾添加await admin.disconnect() - 屏蔽分区器警告:设置环境变量
KAFKAJS_NO_PARTITIONER_WARNING=1即可关闭该警告,或按提示配置LegacyPartitioner
内容的提问来源于stack exchange,提问作者ashique
相关产品推荐
相关产品推荐

