发送失败时需捕获AMQPConnectionError等异常,将消息元数据及base64编码的body结构化为JSON写入带时间戳和ID的本地文件,并调用os.fsync()确保落盘。
发送失败时如何捕获 RabbitMQ 异常并写入本地文件
rabbitmq 的 python 客户端(pika)在连接断开、信道关闭或消息被 nack 时,通常抛出 amqpconnectionerror、channelclosedbybroker 或 unroutableerror 等异常。直接用 try/except 捕获还不够——关键是要把原始消息体(包括 routing_key、exchange、headers、body)完整序列化落盘,否则补发时会丢失上下文。
实操建议:
- 不要只记录错误日志,必须用
json.dumps()将消息元数据 +body(若为 bytes,先 base64 编码)一并写入文件,避免二进制内容损坏 - 文件名带时间戳和唯一 ID,例如
failed_msg_20241105_142309_abc123.json,方便后续按时间扫描 - 写入前加
os.fsync()确保落盘,防止进程崩溃导致日志丢失 - 避免用
logging.FileHandler直接写——它不保证原子写入,也不方便后续程序读取解析
本地失败日志如何结构化存储以便补发识别
补发逻辑需要明确知道:这条日志是发给哪个 exchange?用的什么 routing_key?是否要重设 delivery_mode=2?有没有自定义 headers?这些信息如果只存在日志文本里,解析成本高、易出错。
推荐结构:
{
"timestamp": "2024-11-05T14:23:09.123Z",
"exchange": "order_events",
"routing_key": "order.created",
"properties": {
"delivery_mode": 2,
"headers": {"source": "payment-service"},
"content_type": "application/json"
},
"body_b64": "eyJuYW1lIjogIm9yZGVyLTEyMyJ9"
}
注意点:
立即学习“Python免费学习笔记(深入)”;
-
body_b64字段强制使用 base64,统一处理 str/bytes;补发时再base64.b64decode() -
properties必须包含delivery_mode,否则补发消息可能变成非持久化,重启后丢失 - 不要省略
content_type,消费端依赖它做反序列化判断
补发脚本如何安全地重投失败消息而不重复或漏发
补发不是简单遍历 JSON 文件再调用 channel.basic_publish()。真实场景中,你要考虑:文件是否已被处理过?RabbitMQ 是否已恢复?补发失败要不要二次落盘?
核心策略:
- 补发前先检查 RabbitMQ 连接健康:用
connection.process_data_events(timeout=1)或发一个空心跳帧,避免连不上还硬塞消息 - 每成功重投一条,把对应文件改名为
failed_msg_*.json.done,而不是直接删除——保留痕迹,防误操作 - 补发过程中遇到新异常(如 again
ChannelClosedByBroker),立即中断当前批次,休眠 5 秒后重试连接,不继续刷文件 - 用
glob.glob("failed_msg_*.json")扫描,但按文件名时间戳升序处理,确保老消息优先补发
为什么不能依赖 pika 的自动重连机制来做“优雅补发”
pika 的 BlockingConnection 不支持自动重连后恢复未确认消息;ConnectionParameters(heartbeat=30) 只能维持 TCP 连接存活,无法感知业务层消息是否真正抵达 broker。
典型陷阱:
- 网络抖动时
basic_publish()返回成功,但 broker 实际没收到(TCP ACK 到了,AMQP confirm 没回),此时无异常、无日志、无重试,消息静默丢失 - 启用
publisher_confirms=True后,若 confirm 超时,pika默认抛exceptions.ConnectionClosedByBroker,但不会告诉你哪条消息没 confirm —— 你依然得自己缓存待发消息并比对 - 所有“自动重连 + 重发缓冲区”方案,本质都是在应用层实现可靠投递,RabbitMQ 本身不保证生产者侧的 exactly-once
真正落地时,最易被忽略的是:补发脚本的执行时机。它不该是定时任务轮询,而应作为主服务启动时的初始化步骤 + 独立 CLI 工具双模式存在。否则故障恢复窗口不可控,人工介入成本陡增。


















