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
相关产品推荐
相关产品推荐

