HTTP handler中禁用amqp.Dial,须全局复用RabbitMQ连接;每个请求仅调conn.Channel()获取新channel并立即关闭;publish必须同步启用Confirm模式、设超时500ms并监听ack/nack;消费者需单goroutine处理Ack。

HTTP handler里直接amqp.Dial会炸掉文件描述符
每次请求都调用amqp.Dial,不出几秒就会触发too many open files或connection refused。RabbitMQ连接是长连接,不是HTTP那种短连接,必须全局复用。
启动时一次性建立连接,存到Echo的echo.Context或依赖注入容器里(比如用echov4.WithHTTPErrorHandler配合echo.New().SetHTTPErrorHandler传入全局*amqp.Connection);handler里只调conn.Channel()拿新*amqp.Channel,用完立刻ch.Close()。
- 别把
*amqp.Connection塞进每个request context——它本就是全局的,加锁反而拖慢性能 -
*amqp.Channel绝对不能跨goroutine复用:多个handler并发调ch.Publish会触发fatal error: concurrent map writes - 如果用
echo.Group做模块隔离,也别为每个group建独立连接——还是共用一个*amqp.Connection
ch.Publish必须配ch.Confirm和超时控制
默认ch.Publish是“发了就不管”,broker还没写盘、网络断了、甚至连接已关闭,调用方都收不到任何反馈。要保证消息至少进队列一次,得同步走confirm流程。
正确姿势:在handler内先调ch.Confirm(false)开启确认模式,再ch.Publish,然后用ch.NotifyPublish监听ack/nack,配合context.WithTimeout设500ms上限。超时就记录告警并返回503 Service Unavailable,别让用户以为成功了。
Echo框架 5.1.0 版本源码包下载,适合关注 RealIP 行为变化、StartConfig.Listener、NewDefaultFS 和观测性中间件入口的开发团队。
- 别用
go ch.Publish()丢进后台——handler返回后goroutine可能根本没跑起来 - 别共享同一个
amqp.Publishing实例给多个publish调用:字段如DeliveryMode会被覆盖,导致部分消息非持久化 - 如果publish前要查DB或序列化大结构体,先做完再进confirm流程,避免超时被误判为RabbitMQ问题
消费者端ch.Consume后不能多goroutine直接处理Ack
很多人写msgs := ch.Consume(...); for range msgs { go handle(msg) },结果跑一阵就panic或静默丢消息。因为msg.Ack()必须由声明这个msgs的同一个*amqp.Channel调用,跨goroutine就违反AMQP协议约束。
安全做法只有两种:要么整个消费循环单goroutine串行处理(适合任务耗时稳定、QPS不高的场景);要么用带锁的channel池——但更推荐前者,简单可靠。
- 禁用
autoAck: true,否则消息一到就被删,失败也无法重试 -
deliveryTag只在当前channel生命周期内有效,存到DB或跨服务传递会失效 - 收到
msg.Reject(true)重入队列时,注意避免无限循环——业务层要检查重试次数或加时间戳判断
队列声明必须设durable: true,消息必须设DeliveryMode: amqp.Persistent
这两项漏掉任意一个,RabbitMQ重启后消息全丢。durable是队列属性,Persistent是消息属性,二者是正交的,缺一不可。
声明队列时显式传true:ch.QueueDeclare("task_queue", true, false, false, false, nil);publish时确保amqp.Publishing{DeliveryMode: amqp.Persistent, ...}。别依赖broker默认值——不同版本、不同vhost配置可能不同。
- virtual host要提前创建好,别让代码自动创建(
vhost参数传错会导致连接成功但队列声明失败) - 如果用topic exchange,
routingKey拼写错误不会报错,而是直接丢弃消息——建议在consumer启动时用ch.ExchangeDeclare显式声明exchange并校验 - 测试环境可以关confirm、关durable快速验证逻辑,但CI流水线里必须打开,否则上线即丢消息
defer ch.Close()就能解决的——关键是要接受“每个HTTP请求对应一个临时channel,且必须在返回前完成publish+confirm”这个事实。很多崩溃和丢消息,都源于试图用HTTP思维去套AMQP模型。

















