不能在 HTTP handler 里直接调 ch.Consume(),因为其返回的 Go channel 会持续阻塞推送消息,导致 handler goroutine 永久挂起、连接无法释放、并发耗尽,且连接断开时消费者悄无声息停摆;正确做法是剥离至独立长生命周期 goroutine,配合连接健康检查、自动重建、QoS 控制与幂等确认。

为什么不能在 HTTP handler 里直接调 ch.Consume()
因为 ch.Consume() 返回的是一个 Go channel,它会持续阻塞式地推送消息;如果在 Gin/Echo 的请求 handler 中直接调用,会导致该 goroutine 永久挂起,HTTP 连接无法释放,很快耗尽服务器并发连接数。更严重的是,一旦 RabbitMQ 连接断开或 channel 关闭,ch.Consume() 会返回 error 并退出,但 handler 已经结束,没人监听这个失败——消费者悄无声息地停摆。
正确做法是把消费者逻辑完全剥离到独立的、长生命周期的 goroutine 中启动,且必须配合连接健康检查和自动重建:
- 用
sync.Once初始化全局*amqp.Connection,URL 中含connect_timeout=5和heartbeat=30 - 监听
conn.NotifyClose(),触发后清空旧连接,按指数退避(1s → 2s → 4s)重连 - 每次消费前重新
conn.Channel(),用完立刻ch.Close(),绝不复用 Channel -
ch.Consume()必须在独立 goroutine 中调用,并用for range持续读取消息 channel
delivery.Ack() 放错位置等于丢消息
RabbitMQ 的消费模型是“预取→处理→显式确认”,不是“收到即确认”。很多人把 delivery.Ack(false) 写在业务逻辑中间,比如数据库写入之后、发邮件之前——结果发邮件时 panic,Ack 就被跳过,消息卡在 unack 状态,RabbitMQ 会一直 hold 住它,直到达到 prefetch_count 上限,整个 consumer 停摆。
安全写法只有一条铁律:consumer handler 开头第一行就写 defer delivery.Ack(false),然后做幂等校验和业务处理;只有全部成功才执行 delivery.Ack(true),否则让消息重回队列或进死信队列:
立即学习“go语言免费学习笔记(深入)”;
- 幂等键建议用业务 ID + 事件类型拼接,例如
"order_created_123",用 RedisSETNX校验是否已处理 - 必须设
ch.Qos(1, 0, false)控制预取数量,避免单个 consumer 积压大量 unack 消息 - 不要依赖
delivery.Nack()自动重试——它不保证顺序,且可能无限循环;应主动发到 DLX(死信交换机)并记录日志
订阅主题命名必须带版本和领域前缀
用裸名如 "order" 或 "user" 订阅,看似简单,上线后一升级就炸:v2 版本的订单结构加了字段,老消费者 JSON 反序列化失败 panic;或者新服务误订了旧队列,消费到不该处理的消息。
主题(queue name / routing key / topic)必须体现可演进性:
- 队列名格式:
events.order.created.v1(领域.事件.版本) - Exchange 类型选
topic,绑定时用events.order.*这类通配,方便未来扩展子事件 - 消息体必须是 struct,带
Version string `json:"version"`字段,消费者先校验再解析 - 避免用 fanout exchange 直接广播——它无法按需过滤,所有消费者都得处理每条消息,耦合度反而升高
NATS JetStream 比 RabbitMQ 更适合高频内部事件
如果你的微服务间通信主要是内部状态广播(如配置变更、用户登录踢出、缓存失效),RabbitMQ 的 AMQP 协议和 Exchange/Binding 模型属于过度设计,运维成本高、延迟略高;NATS JetStream 则轻量直接,且原生支持 at-least-once、stream replay、message deduplication。
关键实操点:
- JetStream stream 创建时必须设
Retention: nats.InterestPolicy,否则消息不会持久化 - 用
js.PullSubscribe("ORDERS", nats.Durable("inventory"))实现有状态消费,重启后自动从上次 offset 继续 - 发布时带上
Msg.Header.Set("Nats-Expected-Last-Subject-Sequence", "123")可启用重复检测,天然支持幂等 - 别用
js.PublishAsync()——它不等 broker 确认就返回,网络抖动时消息静默丢失;改用js.Publish()+ context timeout


















