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

如何让Kafkajs Consumer在所有Broker宕机时保持运行或自动恢复?

解决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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 02:45:39