如何让Kafkajs Consumer在所有Broker宕机时保持运行或自动恢复?
问题场景
你有3个运行在9002、9003、9004端口的Kafka Broker,使用KafkaJS创建Consumer后,当所有Broker宕机时,会触发如下错误导致Consumer停止运行:
{"level":"ERROR","timestamp":"2023-03-17T15:07:06.849Z","logger":"kafkajs","message":"[Consumer] Crash: KafkaJSNonRetriableError: Connection error: connect ECONNREFUSED 127.0.0.1:9094","groupId":"groupId","stack":"KafkaJSNonRetriableError: Connection error: connect ECONNREFUSED 127.0.0.1:9094\n at /Users/vn54q4d/projects/path/node_modules/kafkajs/src/retry/index.js:53:18\n at processTicksAndRejections (node:internal/process/task_queues:96:5)"}
以下是两种满足需求的实现方案:
方案1:Broker恢复后自动重启Consumer并继续消费
通过监听Consumer的crash事件,在触发崩溃后实现带退避策略的自动重试逻辑,重新初始化并启动Consumer:
const { Kafka } = require('kafkajs'); const kafka = new Kafka({ clientId: 'your-client-id', brokers: ['localhost:9002', 'localhost:9003', 'localhost:9004'], }); // 定义启动Consumer的复用函数 async function startConsumer() { const consumer = kafka.consumer({ groupId: 'groupId' }); // 监听crash事件,触发后重试启动 consumer.on('crash', async (error) => { console.error('Consumer crashed, retrying...', error); // 指数退避延迟,避免频繁重试 let delay = 1000; const maxDelay = 30000; while (true) { try { await consumer.disconnect(); await startConsumer(); break; } catch (retryError) { console.error('Retry failed, waiting', delay, 'ms...', retryError); await new Promise(resolve => setTimeout(resolve, delay)); delay = Math.min(delay * 2, maxDelay); } } }); try { await consumer.connect(); await consumer.subscribe({ topic: 'your-topic', fromBeginning: false }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { console.log({ value: message.value.toString(), }); }, }); console.log('Consumer started successfully'); } catch (error) { console.error('Failed to start consumer', error); // 首次启动失败直接触发重试逻辑 consumer.emit('crash', error); } } // 启动Consumer startConsumer();
原理:当Consumer因Broker全宕机崩溃时,crash事件被触发,随后通过循环尝试重新连接,每次失败后延迟时间翻倍(最多30秒),直到Broker恢复后成功重启Consumer,继续从上次的偏移量消费数据。
方案2:让Consumer在Broker全宕机时保持运行不停止
修改KafkaJS的重试配置,将连接错误标记为可重试类型,让Consumer自动重试连接而不触发崩溃:
const { Kafka, KafkaJSNonRetriableError } = require('kafkajs'); const kafka = new Kafka({ clientId: 'your-client-id', brokers: ['localhost:9002', 'localhost:9003', 'localhost:9004'], retry: { retries: 999999, // 无限重试(或设置足够大的次数) factor: 2, multiplier: 2, maxRetryTime: 30000, // 自定义重试判断逻辑,将连接错误转为可重试 retryOnError: (error) => { if (error instanceof KafkaJSNonRetriableError && error.message.includes('ECONNREFUSED')) { return true; } // 保留默认的可重试错误判断 return ['RetriableError', 'ConnectionError'].includes(error.name); }, }, }); async function startConsumer() { const consumer = kafka.consumer({ groupId: 'groupId' }); try { await consumer.connect(); await consumer.subscribe({ topic: 'your-topic', fromBeginning: false }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { console.log({ value: message.value.toString(), }); }, }); console.log('Consumer started successfully'); } catch (error) { console.error('Failed to start consumer', error); } } startConsumer();
原理:KafkaJS默认将ECONNREFUSED标记为非可重试错误,会直接触发Consumer崩溃。通过自定义retryOnError函数,将此类连接错误转为可重试,结合高重试次数配置,Consumer会在Broker宕机期间持续尝试重连,不会停止运行,Broker恢复后自动恢复消费。
内容的提问来源于stack exchange,提问作者lal rishav

