如何让KafkaJS消费者无限等待Kafka Broker上线?
KafkaJS消费者无限重试连接配置方案
问题场景
配置KafkaJS消费者后,当Kafka Broker未启动时,客户端仅重试有限次数后停止,无法持续重试直到Broker可用;同时希望Broker掉线后也能自动无限重连。当前代码及错误信息如下:
消费者代码
// Create the kafka client const kafka = new Kafka({ clientId, brokers, }); // Create the consumer const consumer = this.kafka.consumer({ groupId, heartbeatInterval: 3000, sessionTimeout: 30000, }); // Connect the consumer consumer.connect().then(async (res) => { await this.consumer.run({ eachMessage: async ({ topic, partition, message, heartbeat, pause }) => { this.subscriptionRegistrations[topic](topic, partition, message); }, }).catch((err) => { console.log('Error running consumer!', err); }); }).catch(async (err) => { console.log('Error connecting consumer!', err); })
错误信息
Broker未启动时,客户端重试5次后抛出非可重试错误:
{"level":"ERROR","timestamp":"2023-01-27T23:29:58.214Z","logger":"kafkajs","message":"[BrokerPool] Failed to connect to seed broker, trying another broker from the list: Connection error: connect ECONNREFUSED 127.0.0.1:9092","retryCount":0,"retryTime":246} {"level":"ERROR","timestamp":"2023-01-27T23:29:58.463Z","logger":"kafkajs","message":"[Connection] Connection error: connect ECONNREFUSED 127.0.0.1:9092","broker":"localhost:9092","clientId":"CLIENT_ID_TEST","stack":"Error: connect ECONNREFUSED 127.0.0.1:9092\n at TCPConnectWrap.afterConnect [as oncomplete] (node:net:1157:16)"}
最终触发重试次数超限错误:
KafkaJSNonRetriableError Caused by: KafkaJSConnectionError: Connection error: connect ECONNREFUSED 127.0.0.1:9092 at Socket.onError (/path/to/project/node_modules/kafkajs/src/network/connection.js:210:23) at Socket.emit (node:events:526:28) at emitErrorNT (node:internal/streams/destroy:157:8) at emitErrorCloseNT (node:internal/streams/destroy:122:3) at processTicksAndRejections (node:internal/process/task_queues:83:21) [ERROR] 23:29:59 KafkaJSNumberOfRetriesExceeded: Connection error: connect ECONNREFUSED 127.0.0.1:9092
解决方案
KafkaJS默认重试配置有时间和次数限制(默认maxRetryTime为30秒),需显式修改重试参数实现无限重试,同时优化消费者启动逻辑:
1. 修改Kafka客户端重试配置
在创建Kafka实例时,添加retry配置项,设置无限重试规则:
const kafka = new Kafka({ clientId, brokers, retry: { maxRetryTime: Infinity, // 取消重试时间上限,持续重试 initialRetryTime: 1000, // 首次重试间隔1秒 retries: Number.MAX_SAFE_INTEGER, // 最大重试次数设为安全整数上限 factor: 2, // 重试间隔增长因子 multiplier: 1.5, // 每次重试间隔的乘数 }, });
2. 优化消费者启动逻辑
将连接和启动逻辑封装为异步函数,在错误时手动触发重试(作为内部重试的补充保障):
// Create the consumer const consumer = kafka.consumer({ groupId, heartbeatInterval: 3000, sessionTimeout: 30000, }); // 封装消费者启动逻辑 const startConsumer = async () => { try { await consumer.connect(); console.log('Consumer connected successfully'); await consumer.run({ eachMessage: async ({ topic, partition, message, heartbeat, pause }) => { this.subscriptionRegistrations[topic](topic, partition, message); }, }); } catch (err) { console.error('Consumer error, retrying in 5s:', err.message); // 延迟5秒后重试启动 setTimeout(startConsumer, 5000); } }; // 启动消费者 startConsumer();
配置说明
maxRetryTime: Infinity:移除重试的时间限制,让客户端持续尝试连接retries: Number.MAX_SAFE_INTEGER:设置极大的重试次数,近似无限重试initialRetryTime、factor、multiplier:控制重试间隔的初始值和增长速率,避免频繁重试给系统带来压力- 手动重试逻辑:当
connect或run抛出错误时,延迟一段时间后重新启动消费者,作为KafkaJS内部重试的补充
配置完成后,消费者会在Broker未启动时持续重试连接,Broker恢复后自动建立连接并开始消费;当Broker意外掉线时,也会自动触发无限重连,直到Broker恢复。
内容的提问来源于stack exchange,提问作者bapin93
相关产品推荐
相关产品推荐

