重连信号是 readmessage/writemessage 返回 error(如 io.eof、use of closed network connection),而非 panic;需显式检查所有错误,用 conn.notifyclose() 和 ch.notifyclose() 监听关闭,重连逻辑须带 context、退避与 jitter,channel 恢复消费前需 cancel、重新声明 queue 并设置 qos。

ReadMessage/WriteMessage 报错才是重连信号,不是 panic
很多人等 conn.ReadMessage() panic 才去重连,但 Go 的 AMQP 客户端(如 streadway/amqp)和 WebSocket 库一样,**从不 panic**,只返回 error。典型表现是:io.EOF、use of closed network connection、i/o timeout 或 connection reset by peer。
这些错误必须显式检查,不能只写 if err == io.EOF 就退出——因为 io.EOF 只代表对端关闭,而 use of closed network connection 才是本地连接已失效的明确信号。
- 所有
ch.Consume()后的msg.Ack()/msg.Nack()调用前,先确认ch是否仍可用(可通过ch.NotifyClose()监听) -
amqp.Dial()成功后,立刻调用conn.NotifyClose()注册关闭通知 channel,比轮询更及时 - 别在消费 handler 里直接重连——此时 channel 可能还在收消息,应发信号到主 goroutine 统一处理
重连不能裸写 for 循环,必须带上下文与退避
简单 for { conn, err := amqp.Dial(...) if err != nil { time.Sleep(1 * time.Second); continue } } 会吃光 CPU、打爆服务端连接数,且无法响应外部终止指令。
正确做法是把重连逻辑封装进独立 goroutine,并用 context.WithCancel() 控制生命周期:
- 每次重连前先
conn.Close()(如果非 nil),否则旧连接 fd 不释放 - 初始间隔设为
1s,失败后乘以1.8(不是固定 +1),上限封顶60s - 加
jitter:实际休眠时间 =backoff() * (0.9 + rand.Float64()*0.2),防雪崩 - 传
context.WithTimeout(ctx, 5*time.Second)给amqp.Dial(),避免 DNS 卡死拖慢整个退避节奏
channel 断开后如何安全恢复消费,不是简单重启 goroutine
重连成功后,不能直接起新 goroutine 调用 ch.Consume() 就完事。未确认消息可能已在旧 channel 上卡住,新 channel 无法接管它们。
关键动作有三步:
- 调用
ch.Cancel(tag, false)主动取消旧消费者(false表示不 requeue 消息,由业务决定是否重投) - 重新声明 queue(
ch.QueueDeclarePassive()或带passive=false的QueueDeclare()),确保队列还存在 - 重新绑定 exchange 和 routing key;若用的是 auto-delete queue,这步尤其必要
- 恢复 qos 设置:
ch.Qos(1, 0, false),避免一次拉多条导致崩溃后批量丢失
用 taskQuit
多个 goroutine 共享一个 *amqp.Channel 是危险的——它不是并发安全的。ch.Publish() 和 ch.Consume() 同时调用会 panic,且无法保证消息顺序。
推荐结构是「单 channel 单 goroutine」:一个 goroutine 专管消费循环,出错时往 taskQuit channel 发送空结构体;主 goroutine select 监听该 channel,收到后关闭旧 channel、重连、重建新 channel 并启动新消费 goroutine。
这种解耦方式下,你永远只有一组活跃的 conn+ch 实例,不会出现状态混乱或资源泄漏。
最易被忽略的一点:RabbitMQ 的 connection 是线程安全的,但 channel 不是;重连后必须新建 channel,不能复用旧实例——哪怕它还没被 GC 掉。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











