Kafka Raft模式下Node.js服务报主题分区错误排查
问题分析与解决方案
核心问题根源
- 主题创建函数的逻辑缺陷:
- 错误处理中
if (error.code = 3)是赋值操作而非判断,导致非code=3的错误也会触发重试 - 使用
setTimeout递归调用时,外层函数无法等待异步重试完成,会提前返回,导致消费者在主题未完全就绪时就开始订阅
- 错误处理中
- 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
相关产品推荐
相关产品推荐

