必须设autoAck=false,因为autoAck=true时RabbitMQ一发送消息即删除,无论panic、网络中断或业务未执行完,消息必然丢失;须手动调msg.Ack(false)或msg.Nack(true,false),且仅限当前goroutine内调用一次。

为什么不能用 autoAck=true 启动消费者
开 autoAck=true 就等于告诉 RabbitMQ:“你发出去就当我处理完了”,哪怕你的 Go 程序刚收到消息就 panic、网络断了、或者数据库写一半崩了,消息也永远消失了。这不是概率问题,是设计上必然丢失。
必须设为 autoAck=false,并在业务逻辑真正完成后再手动调用 msg.Ack(false)。否则所有“至少一次投递”的保障都形同虚设。
-
msg.Ack(false)只能调用一次,且必须在接收该msg的 goroutine 内执行;跨 goroutine 或 defer 里调用会卡死连接 - 若业务逻辑 panic,需用
defer+recover捕获,并立刻调用msg.Nack(true, false)让消息重入队 - 别漏掉超时控制:用
context.WithTimeout包裹整个处理流程,超时后也应Nack,避免消息长期滞留
goroutine 泄漏和并发失控怎么防
常见错误是每来一条消息就起一个 goroutine:go process(msg)。短时间大量消息进来,会瞬间创建成百上千 goroutine,调度器不堪重负,CPU 和内存飙升,甚至触发 OOM。
正确做法是固定数量的工作 goroutine 复用 channel,类似线程池:
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
立即学习“go语言免费学习笔记(深入)”;
- 启动前确定合理并发数(如
workerCount = runtime.NumCPU()或按压测结果设为 4~8) - 每个 worker goroutine 循环读取同一个
chan amqp.Delivery,不新建不销毁 - 避免在 consumer 循环里直接调用 HTTP handler 或 Beego Controller 方法——它们生命周期不匹配,ACK 会失效
- 若需共享状态(如计数器、缓存),用
sync.Map或带锁结构,别裸写全局 map
消息体过大或序列化错位导致静默失败
Go 消费者收不到消息?可能根本不是连接或 ACK 的问题,而是反序列化失败后直接 panic 或忽略错误,日志里只有一行 invalid character,没人看见。
- 所有
json.Unmarshal/proto.Unmarshal必须检查 error,不能只用_忽略 - 消息体建议控制在 100KB 以内;超大 payload(如图片 base64)应存对象存储,消息只传 URL
- 不同服务间要约定好序列化格式;Go 生产者用
json,Java 消费者却用protobuf,必然解析失败 - 每条消息打日志时带上
msg.MessageId和msg.CorrelationId,否则排查时根本分不清哪条 log 对应哪条消息
重启后 offset 怎么续?本地 channel 不扛这事
用 make(chan, N) 做内存队列,服务一重启,没消费完的消息全丢。生产环境必须依赖外部存储记录消费位置。
- RabbitMQ 场景下,靠
deliveryTag和服务端 ack 状态隐式维护进度,无需额外 offset;但得确保 channel 不被意外关闭 - Kafka 场景下,必须调用
consumer.CommitOffsets(),且不能只靠定时 commit——应在每批成功处理后同步提交 - 若用 Redis Streams,用
XREADGROUP+pending列表,重启时从XPENDING拉未确认消息继续处理 - 别把 offset 存文件——并发写容易损坏;也别用内存变量模拟——进程挂了就归零
Nack 当重试开关,却没配 DLX;或者以为 context.Cancel 能中断正在跑的 goroutine,其实它只影响后续阻塞操作。这些细节,不踩一次很难记住。

















