不能在gin handler里直接调用ch.publish,因http响应写出后连接可能被复用或关闭,此时ch.publish遇网络抖动、channel关闭等会panic或静默丢消息;须用带缓冲channel+独立goroutine异步发送,并每次新建amqp.publishing实例、检查ch.isclosed()、传入context控制超时。

为什么不能在 Gin Handler 里直接调用 ch.Publish
因为 HTTP 响应一旦写出(w.WriteHeader 或 w.Write 返回),底层 net.Conn 就可能被复用或关闭。此时若 ch.Publish 正在执行,遇到网络抖动、RabbitMQ 拒绝连接、channel 已关闭等情况,会触发 panic: send on closed channel,或静默丢消息——这不是异步,是“扔完不管”。
常见错误包括:
-
defer ch.Close()放在 handler 末尾,但 channel 实际已在响应写出后被回收 - 复用已关闭的
amqp.Channel,没检查ch.IsClosed() - 用全局单例
*amqp.Channel并发写入,违反 AMQP 协议(channel 不支持多 goroutine 同时写)
怎么让消息投递不阻塞 Gin 接口响应
核心就一条:Gin handler 只做序列化 + 非阻塞投递;真正发送交给独立 goroutine。
- 声明带缓冲的全局 channel:
var msgCh = make(chan []byte, 1000)(缓冲大小按峰值 QPS × 平均处理延迟预估) - handler 中用
select非阻塞写入:select { case msgCh - 单独 goroutine 消费:
go func() { for payload := range msgCh { publishToRabbitMQ(payload) } }() -
publishToRabbitMQ必须每次新建amqp.Publishing{}实例,不能复用(避免 header map 共享污染) - 必须捕获
amqp.Error(如Code: 404队列不存在),不能直接 panic 或忽略
连接和 channel 生命周期怎么管才靠谱
RabbitMQ 连接是重量级资源,不能每次请求都 amqp.Dial + ch.Close(),否则高并发下会耗尽 TCP 端口或触发 RabbitMQ 连接数上限。
- 连接(
*amqp.Connection)必须全局单例,启动时初始化,用sync.Once保证 - channel 应按业务域隔离(如“订单事件 channel”和“通知 channel”不共用),避免互相干扰
- 必须监听
conn.NotifyClose(),触发自动重连逻辑,否则网络抖动后整个链路静默中断 - 每次发送前检查
ch.IsClosed(),不可用则重建 channel(不是重连 connection) -
deliveryMode: amqp.Persistent必须设为2,否则 broker 重启后消息丢失
消费者手动 ACK 为什么总卡在 unacked 状态
不是 RabbitMQ 不给 ACK,是你根本没发出去。最常见原因是业务逻辑里 return、panic 或未 recover 的 error,导致 msg.Ack(false) 根本没执行。
-
msg.Ack(false)必须放在defer里,且defer前确保msg非 nil - 业务逻辑必须包在
defer func() { if r := recover(); r != nil { msg.Nack(false, true) } }()中 - 不要用
msg.Ack(true)(multiple=true),它会批量确认前面所有消息,容易误确认 - 必须绑定
context.Context到执行层,超时能终止 DB 查询或 HTTP 调用,避免 goroutine 挂死
真正难的不是连上 RabbitMQ,而是让每条消息在连接断、channel 关、panic、超时、队列不存在等几十种异常路径下,依然可追踪、可重试、不丢失、不堆积。这些细节不落地,压测一跑就崩。











