必须用message.newrouter启动,禁用for range手写消费循环;watermill强制以router为唯一入口统一调度消息流,手写循环会导致goroutine泄漏、panic中断、中间件失效,且handler需通过addhandler注册并严格遵循返回值语义控制确认与转发。

必须用 message.NewRouter 启动,禁用 for range 手写消费循环
Watermill 不是“连上 Kafka 就能收消息”的胶水库,它强制以 message.NewRouter 为唯一入口统一调度所有消息流。跳过 Router 直接调 subscriber.Subscribe("topic") + for msg := range messages,上线后必然出现 goroutine 泄漏、panic 全局中断、中间件失效。
-
Router内部用 channel 复用 +sync.WaitGroup管理生命周期,比手动启多个go func() { for range }更轻量可控 - 所有中间件(如
RetryMiddleware、TracingMiddleware)只对router.AddHandler注册的 handler 生效,手写循环完全绕过 - 每个 handler panic 是隔离的;而手写
for range一旦 panic,整个 loop 就退出,下游消息全部卡住
router.AddHandler 的参数和返回值语义必须严格对齐
handler 函数签名固定为 func(*message.Message) ([]*message.Message, error),返回值直接控制消息确认与转发行为,错一个就丢消息或无限重试。
- 返回
nil, nil:仅消费,自动调msg.Ack()(前提是未启用ManualAck) - 返回
[]*message.Message, nil:消费并发布新消息到publishTopic - 返回
nil, err或 panic:触发重试(由中间件控制)或永久卡住(若未配重试) - 启用
ManualAck模式时,必须显式调msg.Ack(),否则 offset 永远不提交
Kafka 场景下 ConsumerGroup 和超时配置必须显式设对
消息卡住、重复消费、反复报 context deadline exceeded,90% 是 Kafka 配置没对齐 broker 策略。
-
consumergroup必须显式指定,否则每个实例都被 Kafka 当作独立消费者,导致所有实例都收到同一条消息 - handler 执行时间必须小于 Kafka 的
session.timeout.ms(默认 45s),否则触发 rebalance,消费暂停 - 推荐在 handler 内包一层
ctx, cancel := context.WithTimeout(msg.Context(), 30*time.Second),并确保该 timeout 小于 broker 配置 - 别依赖 Publisher 的重发机制防丢——
kafka.NewPublisher只是发送客户端,和数据库事务完全无关
防丢消息必须走 outbox.NewPublisher,不能靠配置硬扛
订单、支付等关键事件“发出去但 DB 没写成功”,或“DB 写成功但消息没发出去”,本质是事务与消息发布没绑定。Watermill 不提供原子性保障,得自己加 Outbox。
- 本地开发可用
gochannel.NewPublisher,但生产环境必须切到kafka.NewPublisher或amqp.NewPublisher - 真正防丢要引入
outbox.NewPublisher(db, publisher, logger),把事件先写进数据库 outbox 表,再由轮询协程发出 -
outboxPub.Publish(ctx, "topic", msg)这行必须在tx.Commit()前调用,才能保证和业务操作在同一个事务里 - Outbox 表结构需包含
payload、topic、published_at字段,轮询逻辑需幂等处理已发消息
动态路由发现、服务注册集成这些高级能力,都建立在 Router 正确初始化和 handler 语义对齐的基础上。没跑通基础消费模型之前,别碰 etcd/Consul 自动注册——那只会把问题藏得更深。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











