因为amqp.Publish默认为“fire-and-forget”模式,不启用publisher confirm机制,不校验Broker是否接收成功,也不等待ack/nack,导致高并发或网络抖动时消息静默丢失。

为什么直接用 amqp.Publish 会丢消息
很多团队一上来就用 channel.Publish 发消息,结果在高并发或网络抖动时发现事件“消失”了——根本没进队列。这不是代码写错了,而是 RabbitMQ 默认开启的是“fire-and-forget”模式,amqp.Publish 不校验 broker 是否真正接收成功,也不等 confirm。
真正可靠广播的前提是:启用 publisher confirms(不是事务),并配合重试 + 唯一消息 ID(比如用 correlation_id 或自定义 message_id 字段)。否则哪怕服务 A 认为自己发出去了,RabbitMQ 可能因内存满、镜像同步延迟等原因静默丢弃。
- 必须调用
channel.Confirm开启确认模式,且在publish后监听channel.NotifyPublish的 ack/nack - 不要依赖
amqp.Table里的delivery_mode: 2(持久化)来代替 confirm——它只保证入磁盘,不保证路由成功 - nack 时别直接 panic,应记录日志 + 放入本地重试队列(例如用
time.AfterFunc延迟重发,避免雪崩)
如何让多个微服务订阅同一事件但互不干扰
常见误区是所有服务都绑定到同一个 queue,结果变成竞争消费(只有 1 个服务收到),而不是广播。要实现真正的“事件广播”,必须用 exchange + fanout + 独立 queue。
每个服务启动时,应声明一个**排他性、自动删除、非持久化**的 queue,并绑定到同一个 fanout exchange。这样每条消息会被复制到每个 queue,各服务独立消费,互不影响。
立即学习“go语言免费学习笔记(深入)”;
- queue 名字建议包含服务名和实例标识(如
"svc-order-7f9a"),避免不同环境冲突 - 不要设
durable: true——排他 queue 本身就不该持久化;持久化 queue + fanout 容易堆积,且重启后残留旧消息干扰新实例 - 务必设置
autoDelete: true,否则服务异常退出后 queue 残留,下次启动可能重复绑定失败
怎么避免事件重复消费导致状态错乱
RabbitMQ 不保证 exactly-once,只提供 at-least-once。网络分区、消费者 crash 后 reconnect、手动 nack 都可能导致同一条消息被投递多次。靠 ACK 机制本身无法解决,必须业务层兜底。
核心做法是:所有事件消息带全局唯一 message_id(推荐用 UUIDv4),消费者处理前先查本地幂等表(如 Redis 或 DB 表 event_processed),存在则跳过。
-
message_id必须由发布方生成并写入 AMQP header(amqp.Publishing.Headers),不能靠 broker 自动生成——RabbitMQ 的message_id字段默认为空,且不强制填充 - 幂等检查要在
delivery.Ack()之前完成,否则可能漏判;失败时用delivery.Nack(requeue: false)丢弃,别 requeue - Redis key 建议用
"evt:<code>message_id" + TTL(如 24h),避免无限膨胀
Go 客户端连接崩溃后如何自动恢复
原生 amqp.Connection 和 amqp.Channel 都不支持自动重连。一旦网络闪断或 broker 重启,connection.NotifyClose 会触发,但 channel 已失效,继续 publish 会 panic。
不能简单地在 notify 后 new connection —— 多 goroutine 并发重建容易冲突,且未处理完的 pending publish 会丢失。正确做法是封装一层带状态机的连接管理器,控制重建节奏。
- 用
sync.RWMutex保护当前 channel,每次 publish 前RLock,失败则RLock→Unlock→Lock→ 重建 - 重建 channel 时重新声明 exchange/queue/bindings,不要复用旧声明——旧 channel 的资源已不可用
- 加指数退避(如 100ms → 200ms → 400ms),避免瞬间重连风暴打挂 broker
最麻烦的其实是消费者端:reconnect 后要重新 consume,但旧的 delivery channel 可能还在读,得用 context.WithCancel 主动关闭旧 goroutine,再启新的。


















