关于RdKafka::Producer在Broker重启时丢失队列消息的问题咨询
嗨,我来帮你梳理下这个问题的解决方案~
首先明确你的核心场景:你用librdkafka v2.1.1开发的C++生产者,在Kafka Broker(3.4.0版本,运行在WSL2 Ubuntu 22.04.1)短暂停机重启后,出现了离线缓存消息丢失的情况;但把Broker的auto.create.topics.enable设为true就没有问题。你的生产者已经配置了最大重试次数(retries=2147483647)和超长投递超时(delivery.timeout.ms=86400000),主题也是提前创建好的(1副本1分区),问题大概率出在Broker重启后的元数据同步不及时,以及librdkafka对主题存在性的错误判断上。
为什么auto.create.topics.enable=true能解决问题?
当Broker重启后,生产者本地缓存的元数据还没更新,此时会误判主题不存在。如果auto.create开启,生产者会向Broker发送创建主题的请求,而Broker此时已经加载了原有主题的元数据,会返回“主题已存在”的响应,顺便把最新的主题/分区元数据同步给生产者,缓存的消息就能正常重试投递了。但当auto.create关闭时,生产者收到“主题不存在”的错误后,可能会直接判定这些消息属于不可重试错误,最终导致消息丢失。
针对auto.create.topics.enable=false的解决方案
你可以试试下面几个调整方向:
缩短元数据刷新间隔
调整librdkafka的元数据相关配置,让生产者更频繁地主动刷新元数据,这样Broker重启后能更快感知到主题的存在:- 设置
metadata.max.age.ms=30000(默认是300000,即5分钟),控制全局元数据的最大过期时间,到期就会触发刷新 - 设置
topic.metadata.refresh.interval.ms=10000(默认是3600000,即1小时),针对单个主题的元数据刷新间隔,调小后能更快获取主题的最新状态
- 设置
在回调中手动触发元数据刷新
在你的Delivery Report Callback里,当捕获到"Broker: Unknown topic or partition"这类错误时,可以手动触发一次元数据刷新,强制生产者去Broker拉取最新的元数据:// 示例:在回调中触发指定主题的元数据刷新 if (rkt && strcmp(errstr, "Broker: Unknown topic or partition") == 0) { rd_kafka_metadata_t *metadata; // 拉取指定主题的元数据,超时时间设为5秒 rd_kafka_metadata(rk, 0, rkt, &metadata, 5000); rd_kafka_metadata_destroy(metadata); }这样能让生产者更快确认主题已经恢复,进而重新投递缓存的消息。
检查Broker重启后的主题状态
偶尔Broker重启后,可能因为log目录权限、磁盘问题导致主题无法正常加载,此时即使提前创建过主题,Broker也会返回“主题不存在”。你可以在Broker重启后,用Kafka自带工具确认主题状态:kafka-topics.sh --bootstrap-server <你的Broker地址> --describe --topic <你的主题名>确保输出里显示主题存在,分区状态为
Leader: 0(或你的Broker ID),没有异常信息。优化重试逻辑配置
确保retry.backoff.ms(默认100)和retry.backoff.max.ms(默认1000)的配置合理,让生产者在遇到错误时不会过快放弃重试,而是有足够的时间等待Broker恢复。
你可以先从调整元数据刷新间隔开始尝试,这个操作最简单,应该能解决大部分场景下的问题。如果还是不行,可以再检查Broker的主题加载状态,或者在回调里手动处理元数据刷新。
备注:内容来源于stack exchange,提问作者ed_soft

