watermill不是开箱即用事件总线,必须手动组装router、显式注册handler并配置序列化器;router为一次性对象不可复用,需在main中单例启动,topic与marshaler必须严格匹配,kafka需固定group.id以保障消费语义。

Watermill 在 Go 项目里不是“开箱即用”的事件总线,它不提供全局单例或自动注册机制;你得自己组装消息路由器、定义序列化方式、显式绑定 handler,否则 PubSub 发出去的消息根本没人消费。
为什么 Router 必须手动启动且不能复用
Watermill 的 Router 是一次性对象:启动后状态不可重置,内部维护 goroutine 和 channel 生命周期。如果在测试中反复创建/启动/停止,容易触发 panic: send on closed channel 或 goroutine 泄漏。
- 每个服务实例只应创建并启动一个
Router,通常放在main()或应用初始化阶段 - 不要把
Router当作依赖注入到多个 handler 中再分别调用Run—— 它本身已负责调度所有注册的 handler - 若需隔离环境(如测试),用
watermill.NewMemoryMessageRouter替代 Kafka/NATS 实现,避免外部依赖干扰
Handler 注册时必须指定唯一 Topic 和明确的 Unmarshaler
Watermill 不自动推断消息结构。如果你用 JSON 序列化但没配 JSONMarshaler,或 Topic 名拼写与发布端不一致,消息就会静默丢弃 —— 不报错,也不进 handler。
在 Go 中使用 google/wire 实现编译时依赖注入——wire.NewSet、wire.Build、wire.Bind(接口→实现)、wire.Struct、wire.Value、wire.Interface
- 发布端和订阅端的
Topic字符串必须完全一致(包括大小写、空格、下划线) - 始终显式设置
router.AddHandler的第 4 个参数为watermill.DefaultJSONMarshaler{},除非你自定义了二进制协议 - handler 函数签名必须是
func(msg *message.Message) error,返回nil表示成功,非nil会触发重试(默认 3 次)
Kafka 作为 PubSub 时,group.id 决定消费行为
Watermill 对 Kafka 的封装基于 sarama,但隐藏了 consumer group 管理细节。如果你没在 KafkaConfig 中设 GroupID,它会用随机字符串,导致每次重启都从头消费,无法实现 “至少一次” 语义。
-
GroupID必须固定且业务相关(如"order-service-processor"),否则 offset 不会持久化 - 生产环境务必配置
OffsetsInitial:设为sarama.OffsetNewest避免回溯历史消息,或sarama.OffsetOldest仅用于首次全量重建 - 注意 Kafka broker 版本兼容性 —— Watermill v1.5+ 要求 broker ≥ 2.0,低于此版本可能卡在 metadata 请求
最常被跳过的一步是验证 Message.Metadata 是否携带必要上下文(比如 trace-id、source-service)。Watermill 默认不透传元数据,需要在 publisher 侧手动塞入,consumer 侧再提取 —— 这部分逻辑不在框架内,得你自己补。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










