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

Kafka Raft模式下Node.js服务报主题分区错误排查

问题分析与解决方案

核心问题根源

  1. 主题创建函数的逻辑缺陷:
    • 错误处理中if (error.code = 3)是赋值操作而非判断,导致非code=3的错误也会触发重试
    • 使用setTimeout递归调用时,外层函数无法等待异步重试完成,会提前返回,导致消费者在主题未完全就绪时就开始订阅
  2. KRaft模式的元数据同步延迟:在KRaft集群中,主题创建后需要一定时间完成元数据同步到所有controller和broker,直接后续订阅会因元数据不一致触发UNKNOWN_TOPIC_OR_PARTITION错误

修复后的主题创建函数

const delay = (ms) => new Promise(resolve => setTimeout(resolve, ms));

const createTopicsAndPartitions = async (topicList, kafka, topics = null) => {
  let admin;
  topics = topics ? topics : topicList.map((x) => ({
    topic: x,
    numPartitions: x === "sem.initialize.database" ? 1 : 3,
    replicationFactor: 3
  }));

  try {
    admin = kafka.admin();
    await admin.connect();

    // 尝试创建主题,等待分区leader选举完成
    await admin.createTopics({
      topics: topics,
      waitForLeaders: true
    });

    // 轮询确认主题元数据同步完成
    await Promise.all(topics.map(async (topicConfig) => {
      const topic = topicConfig.topic;
      // 最多尝试5次,每次间隔1秒
      for (let i = 0; i < 5; i++) {
        const topicMetadata = await admin.fetchTopicMetadata({ topics: [topic] });
        const partitions = topicMetadata.topics[0]?.partitions;
        if (partitions && partitions.every(p => p.leader !== -1)) {
          return;
        }
        await delay(1000);
      }
      throw new Error(`主题 ${topic} 长时间未就绪`);
    }));

  } catch (error) {
    console.trace(error);
    // 仅处理UNKNOWN_TOPIC_OR_PARTITION错误,限制重试次数避免死循环
    if (error.code === 3 && (!error.retryCount || error.retryCount < 3)) {
      await delay(5000);
      await createTopicsAndPartitions(null, kafka, topics.map(t => ({ ...t, retryCount: (t.retryCount || 0) + 1 })));
    } else {
      throw error;
    }
  } finally {
    if (admin) {
      await admin.disconnect();
    }
  }
};

消费者启动流程优化

const startConsumer = async () => {
  let topicList = ["not_important"];
  if (config.IS_CLOUD) {
    topicList = [...topicList, "sem.initialize.database"];
  }

  // 等待主题完全创建并就绪
  await createTopicsAndPartitions(topicList, kafka);
  
  try {
    await consumer.connect();
    
    // 订阅前先验证主题元数据
    const admin = kafka.admin();
    await admin.connect();
    await admin.fetchTopicMetadata({ topics: topicList });
    await admin.disconnect();

    await consumer.subscribe({
      topics: topicList,
      fromBeginning: false,
    });

    await consumer.run({
      eachMessage: async ({ topic, partition, message }) => {
        try {
          switch (topic) {
            case "sem.initialize.database":
              return await initializeDatabase(message);
          }
        } catch (error) {
          console.trace(error);
        }
      },
      // 监听元数据错误,自动重启消费者
      onError: async (error) => {
        if (error.code === 3) {
          console.warn("检测到主题元数据错误,重启消费者");
          process.emit("RESTART_CONSUMER");
        }
      }
    });
  } catch (error) {
    console.trace("消费者启动失败:", error);
    setTimeout(() => startConsumer(), 5000);
  }

  // 保留原有的重启、关闭逻辑
  const restart = async () => {
    try {
      await consumer.stop();
      await startConsumer();
      logger.info("Consumer has been restarted.");
    } catch (error) {
      logger.error(console.trace(error));
    }
  };
  const shutdown = async () => {
    logger.info("Shutting down the Kafka consumer...");
    try {
      await consumer.disconnect();
      logger.info("Kafka consumer disconnected");
      process.exit(0);
    } catch (error) {
      console.error("Error while disconnecting Kafka consumer:", error);
      process.exit(1);
    }
  };
  process.on("SIGINT", shutdown);
  process.on("SIGTERM", shutdown);
  process.on("RESTART_CONSUMER", restart);
};

关键修复点说明

  • waitForLeaders: true:确保KafkaJS等待主题所有分区的leader选举完成后再返回
  • 元数据轮询:主动检查主题分区的leader状态,确认元数据同步到所有节点
  • 修复错误判断逻辑:将赋值操作改为严格相等判断,避免错误触发重试
  • 异步重试等待:用await delay()替代setTimeout,确保主题未就绪时不会提前启动消费者
  • 消费者错误监听:遇到元数据错误时自动重启,避免持续丢失消息

内容的提问来源于stack exchange,提问作者Furkan YIlmaZ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 03:24:54