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

Node.js Kafkajs消费者连接Kafka集群失败,如何匹配Spring配置?

问题描述

我使用Kafkajs库编写了Kafka消费者代码,但运行时出现连接错误,而Spring Kafka消费者可以正常消费同一个Topic。

我的Kafkajs代码

const { Kafka } = require('kafkajs');

// Create a Kafka consumer
const kafka = new Kafka({
  clientId: 'my-kafka-consumer',
  brokers: ['broker-address:9092'],
  ssl: true,
  sasl: {
     mechanism: 'plain',
     username: 'api-key',
     password: 'api-secret'
  }
});

const consumer = kafka.consumer({ groupId: 'test-group' }); // Specify your consumer group ID

// Function to consume messages from Kafka topic
const consumeMessages = async () => {
  try {
    await consumer.connect();
    await consumer.subscribe({ topic: 'test topic' }); 
    await consumer.run({
      eachMessage: async ({ topic, partition, message }) => {
        console.log({
          topic,
          partition,
          offset: message.offset,
          value: message.value.toString()
        });
      }
    });
  } catch (error) {
    console.error('Error consuming messages:', error);
  }
};

consumeMessages();

报错信息

{"level":"ERROR","timestamp":"2024-04-19T11:10:41.217Z","logger":"kafkajs","message":"[Connection] Connection error: read ECONNRESET","broker":"broker-address:9092","clientId":"my-kafka-consumer","stack":"Error: read ECONNRESET\n    at TCP.onStreamRead (node:internal/stream_base_commons:217:20)"}
{"level":"ERROR","timestamp":"2024-04-19T11:10:41.221Z","logger":"kafkajs","message":"[BrokerPool] Failed to connect to seed broker, trying another broker from the list: Connection error: read ECONNRESET","retryCount":0,"retryTime":262}

可用的Spring Kafka配置(能正常连接)

spring:
  kafka:
    bootstrap-servers: server-url:9092
    properties:
      security.protocol: SASL_SSL
      basic.auth.credentials.source: USER_INFO
      sasl:
        mechanism: PLAIN
        jaas.config: org.apache.kafka.common.security.plain.PlainLoginModule required username='${CLUSTER_API_KEY}' password='${CLUSTER_API_SECRET}';
      schema.registry:
        url: <schema url>
        basic.auth.user.info: ${SCHEMA_REGISTRY_KEY}:${SCHEMA_REGISTRY_SECRET}
      auto.register.schemas: false

请问如何调整Kafkajs消费者的配置参数,匹配Spring配置以成功连接到Kafka集群?


解决方案

要匹配Spring Kafka的配置,需调整Kafkajs的以下核心配置项:

1. 修正Broker地址

Spring配置中bootstrap-servers为server-url:9092,但你的Kafkajs代码使用的是broker-address:9092,先确保两者地址完全一致,这是连接失败最常见的原因。

2. 明确指定安全协议

Spring配置的security.protocol: SASL_SSL对应Kafkajs中的securityProtocol: 'sasl_ssl',需显式配置以确保安全机制对齐,避免隐式参数不匹配。

3. 校验SASL认证信息

Spring使用的PLAIN机制与Kafkajs配置一致,但需确保username和password对应Spring中的CLUSTER_API_KEY和CLUSTER_API_SECRET,而非占位符值。

4. SSL额外配置(按需添加)

若集群要求验证SSL证书,可补充SSL相关配置:

ssl: {
  rejectUnauthorized: true,
  ca: [fs.readFileSync('/path/to/ca-cert.pem', 'utf-8')]
}

5. Schema Registry适配(消费Avro消息时)

Spring配置包含Schema Registry,但Kafkajs本身不处理Schema,若消费Avro序列化消息,需使用@kafkajs/confluent-schema-registry库处理:

const { SchemaRegistry } = require('@kafkajs/confluent-schema-registry')
const registry = new SchemaRegistry({
  host: '<schema url>',
  auth: {
    username: 'SCHEMA_REGISTRY_KEY',
    password: 'SCHEMA_REGISTRY_SECRET'
  }
})

调整后的完整Kafkajs代码

const { Kafka } = require('kafkajs');
// 消费Avro消息时引入Schema Registry库
// const { SchemaRegistry } = require('@kafkajs/confluent-schema-registry');

const kafka = new Kafka({
  clientId: 'my-kafka-consumer',
  // 替换为Spring配置中的server-url:9092
  brokers: ['server-url:9092'],
  // 明确指定安全协议
  securityProtocol: 'sasl_ssl',
  ssl: true,
  sasl: {
    mechanism: 'plain',
    // 替换为实际的集群API密钥和密钥
    username: 'your-cluster-api-key',
    password: 'your-cluster-api-secret'
  }
});

// 消费Avro消息时初始化Schema Registry
// const registry = new SchemaRegistry({
//   host: '<schema url>',
//   auth: {
//     username: 'your-schema-registry-key',
//     password: 'your-schema-registry-secret'
//   }
// });

const consumer = kafka.consumer({ groupId: 'test-group' });

const consumeMessages = async () => {
  try {
    await consumer.connect();
    // 确保Topic名称与Spring消费的一致,避免拼写或空格错误
    await consumer.subscribe({ topic: 'test topic' });
    await consumer.run({
      eachMessage: async ({ topic, partition, message }) => {
        // 消费Avro消息时反序列化
        // const decodedValue = await registry.decode(message.value);
        console.log({
          topic,
          partition,
          offset: message.offset,
          value: message.value.toString()
          // value: decodedValue
        });
      }
    });
    console.log('消费者已成功连接并开始监听Topic');
  } catch (error) {
    console.error('消费消息出错:', error);
    // 可选:添加重连逻辑
    setTimeout(consumeMessages, 5000);
  }
};

consumeMessages();

额外排查点

  • 用telnet server-url 9092或nc -zv server-url 9092测试Broker端口的网络连通性;
  • 确认消费者组ID无配置冲突;
  • 查看Kafka集群日志,排查是否存在SASL认证失败的详细信息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 21:24:54