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

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等权限
  • 修复Admin客户端使用:消费者代码中listTopics函数连接Admin后未断开,建议在函数末尾添加await admin.disconnect()
  • 屏蔽分区器警告:设置环境变量KAFKAJS_NO_PARTITIONER_WARNING=1即可关闭该警告,或按提示配置LegacyPartitioner

内容的提问来源于stack exchange,提问作者ashique

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 14:24:51