streadway/amqp 库需手动通过 conn.notifyclose 实现重连,因其能捕获tcp断连、broker强制关闭等全量异常,而 ch.notifyclose 仅响应channel级错误;重连须重建连接、channel及所有声明资源,并用缓冲channel避免阻塞,配合指数退避与显式关闭旧连接防泄漏。

streadway/amqp 库本身不提供自动重连,必须手动监听 NotifyClose 并重建连接和 Channel,否则连接断开后程序会 panic 或静默失败。
为什么 NotifyClose 是重连唯一可靠入口
amqp.Connection 和 amqp.Channel 都提供 NotifyClose 方法,但只有 conn.NotifyClose 能捕获底层 TCP 断连、Broker 强制关闭(如 code=320)、认证失效等全量异常;ch.NotifyClose 只反映 Channel 级错误(如队列不存在),无法感知连接已死。
- Broker 重启、网络闪断、防火墙中断都会触发
conn.NotifyClose发送*amqp.Error - 若只监听
ch.NotifyClose,连接已断但 Channel 未显式关闭时,后续ch.Publish会直接 panic:invalid memory address or nil pointer dereference - 必须用带缓冲的 channel(如
make(chan *amqp.Error, 1)),否则协程可能阻塞在发送端
重连时必须重建 Channel 和所有声明资源
Connection 重建后,旧的 *amqp.Channel 对象彻底失效,所有队列、交换机、绑定关系都需重新声明 —— RabbitMQ 不保留这些元数据在新连接上。
- 不能复用旧
ch:调用ch.Publish会返回Exception (504) Channel is closed - 必须按顺序重做:
conn.Channel()→ch.ExchangeDeclare()→ch.QueueDeclare()→ch.QueueBind() - 消费者需重新调用
ch.Consume(),且注意autoAck=false时未确认消息会重新入队,可能重复消费 - 若使用
channel.NotifyPublish或channel.NotifyReturn,也要在新 Channel 上重新注册
避免无限递归重连和 Goroutine 泄漏
裸写 reopen() 递归调用或无节制启 goroutine 容易导致进程夯住或 OOM。
- 不要用
defer conn.Close()在重连函数里 —— 新连接成功前旧连接可能还活着,defer会堆积未执行的 Close - 每次重连前显式检查并关闭旧连接:
if r.conn != nil { r.conn.Close() },否则文件描述符泄漏 - 用指数退避(如
time.Sleep(time.Second )代替固定延时,防止雪崩式重连请求打满 Broker - 重连失败时记录日志并退出 goroutine,不要
continue空转消耗 CPU
生产环境建议直接换封装库
自己实现健壮重连要处理的边界太多:连接风暴抑制、Channel 恢复幂等性、消息暂存重发、消费者位点回溯……
-
wagslane/go-rabbitmq封装了amqp091-go,内置重连管理器和 backoff 策略,NewConn返回的*Conn自动处理底层连接生命周期 - 若必须用
streadway/amqp,至少把重连逻辑抽成独立结构体,带状态机(Connecting/Connected/Failed)和最大重试次数限制 - 关键点:重连不是“连上了就行”,而是“连上 + 所有业务 Channel 和消费者全部就绪”才算成功,中间任何一步失败都应降级或告警
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











