必须用outbox模式将事件写入数据库表再异步投递,禁用http handler直调publisher.publish;router需服务启动时全局初始化并长期运行,handler须用middleware.recoverer隔离panic,outbox轮询协程生命周期独立于http请求。

直接用 Echo 的 HTTP handler 调 publisher.Publish 发事件,看似简单,但上线后大概率出现消息丢失、goroutine 泄漏、HTTP 请求超时却消息未发出——根本原因在于没把 Watermill 的生命周期、错误隔离和事务语义对齐到 Web 服务模型里。
HTTP handler 里不能直接调 publisher.Publish
Echo 的 handler 是短生命周期、阻塞式执行的,而 Watermill 的 publisher.Publish 在 Kafka/RabbitMQ 场景下是异步投递(不等 broker 确认就返回),且无内置重试或 fallback。常见后果:
- 网络抖动时
Publish返回nil,但消息实际未送达,handler 却已返回 200 - 若用
gochannel.NewPublisher测试,消息只存在内存 channel 中,进程重启即丢 - 没包
context.WithTimeout,Kafka 发送卡住会拖垮整个 HTTP 连接池
正确做法是:handler 只负责“记录意图”,把事件写入 outbox 表(与业务 DB 同事务),再由独立 goroutine 异步投递。示例:
// handler 中
err := db.Transaction(func(tx *gorm.DB) error {
if err := tx.Create(&Order{...}).Error; err != nil {
return err
}
// 必须在 tx.Commit() 前调用
return outboxPub.Publish(ctx, "order_topic", msg)
})
Router 必须在服务启动时初始化并长期运行
Echo 启动后,message.NewRouter 实例必须常驻内存并 go router.Run(),不能每次 HTTP 请求都新建一个。否则:
- 每个新 Router 都会启一堆 goroutine 消费消息,实例多时迅速突破 5000 goroutine 上限
- 中间件(如
middleware.Retry、middleware.Recoverer)只在 Router 启动时注册一次,手建 Router 就失效 - Subscriber 的
ConsumerGroup无法复用,Kafka 认为是新消费者,触发 rebalance 导致重复消费
推荐结构:在 main.go 初始化 Router,并传给 Echo 的 echo.Context 或全局变量;Handler 内只取用,不创建。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
Subscriber handler panic 会中断整个 Router,必须显式 recover
Watermill 的 Router 默认不包裹 recover(),一旦某个 topic 的 handler panic,整个 router.Run() 退出,所有 topic 停止消费。这不是 bug,是设计选择——它强制你主动处理错误边界。
- 别依赖 Echo 的全局 panic middleware,它对 Router 内部 goroutine 无效
- 必须用
middleware.Recoverer包裹每个AddHandler,且日志要带msg.UUID方便追踪 - handler 函数内仍需对 DB 查询、HTTP 调用等外部依赖加
context.WithTimeout,避免卡死
示例注册方式:
router.AddHandler(
"order-processor",
"order_topic",
subscriber,
"processed_topic",
publisher,
middleware.Recoverer(logger),
)
Outbox 轮询协程不能和 Echo 共享 context
Outbox 模式需要一个常驻 goroutine 轮询数据库表并调 publisher.Publish。这个 goroutine 的生命周期必须独立于 HTTP 请求 context:
- 若用
echo.Context.Request().Context()启动轮询,HTTP 请求结束时 context cancel,轮询立即退出 - 轮询应使用
context.Background()+ 自定义 cancel func,在服务 shutdown 时统一停止 - 轮询间隔建议 100–500ms,太短加重 DB 压力,太长增加事件延迟
真正难的不是写轮询逻辑,而是确保它和 Router、DB 连接池、HTTP server 的 shutdown 顺序一致:先停轮询,再停 Router,最后关 DB 连接池和 Echo server。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










