不能在Gin handler里直接调ch.Consume(),因其返回的<-chan amqp.Delivery会阻塞并持续接收消息,而Gin handler必须毫秒级返回HTTP响应,否则导致连接卡死、请求挂起。

为什么不能在Gin handler里直接调ch.Consume()
因为ch.Consume()返回的是一个<-chan amqp.Delivery,它会阻塞并持续接收消息——而Gin handler必须在毫秒级内返回HTTP响应。一旦你把它写进路由函数,整个HTTP连接就卡死在for-range循环里,后续请求全被挂起。
常见错误写法:msgs, _ := ch.Consume(...); for msg := range msgs { ... },这会让单个HTTP goroutine变成消费者goroutine,既无法并发,又违反了HTTP生命周期约束。
- 正确姿势:消费逻辑必须在
main()或init()阶段启动独立goroutine,与HTTP服务并行运行 - 不要把
ch.Consume()放在任何gin.HandlerFunc内部 - 每个消费者goroutine应绑定专属
amqp.Channel,避免多个goroutine共用一个channel导致PRECONDITION_FAILED - unknown delivery tag
autoAck: false不是可选项,是保底开关
设成true等于告诉RabbitMQ:“我收到了,删吧”,哪怕你的业务逻辑还没执行第一行代码、甚至刚解包完JSON就panic,消息也永远消失了。线上出问题时90%的“消息丢失”都源于此。
必须显式传false:ch.Consume(queueName, "", false, false, false, false, nil),然后在业务逻辑成功后手动调msg.Ack(false)。
- 别在
for range msgs外写defer msg.Ack()——msg是循环变量,defer会捕获到最后一次迭代的值,前面所有消息都漏确认 - 失败时用
msg.Nack(false, true)(requeue=true),让消息回到队列头部重试,这是防丢底线 - 处理前先拷贝body:
data := append([]byte(nil), msg.Body),否则后续迭代可能覆盖内存,导致JSON解析错乱
消费者goroutine怎么扛住连接断开?
RabbitMQ连接随时可能因网络抖动、broker重启或心跳超时断开。裸写的for range msgs一断就退出,再不恢复,整个消费链路静默中断。
必须监听conn.NotifyClose()和ch.NotifyClose(),并在关闭事件触发后重建连接和channel,再重新调用ch.Consume()。
- 不要用全局单例
*amqp.Channel长期持有——channel不可复用,每次断连后必须新建 - 重连逻辑要带指数退避(如1s→2s→4s),避免雪崩式重连打爆broker
- 在
ch.Consume()前检查ch.IsClosed(),已关闭则跳过直接重建 - consumer goroutine本身要包裹在
for无限循环里,确保异常退出后能重启
消息体反序列化容易踩哪些坑?
Go的json.Unmarshal对零值、nil map、嵌套结构极其敏感。生产者发了个params: nil,消费者用struct{ Params map[string]string }接收,就会panic而不是跳过。
定义任务结构体时必须明确控制序列化行为:
- 所有字段加
json:"field_name,omitempty",避免意外透出零值污染下游 - 关键字段(如
ID、Type)不加omitempty,且Unmarshal后立刻校验非空 - 拒绝
interface{}泛型接收,按ContentType或Typeheader路由到具体子类型(如SendEmailTask、ProcessImageTask) - 消息必须设
DeliveryMode: amqp.Persistent,否则broker重启后未消费消息全丢
最麻烦的其实是deliveryTag——它只是当前channel内单调递增的整数,跨goroutine或复用channel确认会直接触发PRECONDITION_FAILED,这点连很多老手都会忽略。


















