Watermill 必须通过 Router 启动,不可直接 for range 消费;Router 是调度中枢,负责 channel 复用、生命周期管理、panic 隔离和中间件生效;正确做法是使用 message.NewRouter 并调用 AddHandler;Kafka 卡顿或重复多因 ConsumerGroup 或 Ack 配置错误;防消息丢失需结合 Outbox 模式而非仅依赖 Publisher;Google Cloud Pub/Sub 集成需正确配置 ProjectID 和唯一订阅名。

Watermill 必须用 Router 启动,不能直接 for range 消费
跳过 message.NewRouter、手写 for msg := range messages 是上线后 goroutine 泄漏和 panic 全局中断的头号原因。Router 不是可选包装,而是 Watermill 的调度中枢:它复用 channel、用 sync.WaitGroup 管理生命周期、隔离每个 handler 的 panic,还让中间件(如 RetryMiddleware)生效。
常见错误现象:
- 本地跑 3 个 topic 就启了 3 个独立 goroutine,线上扩到 20+ topic 后内存持续上涨
- 一个 handler panic,整个 subscriber loop 停摆,其他 topic 消息全部积压
- 加了日志中间件,但手写循环里完全没打印——因为中间件只对
router.AddHandler生效
正确做法:
- 初始化
router := message.NewRouter(message.RouterConfig{}) - 每个 topic 调用
router.AddHandler("order.created", "order_topic", subscriber, publisher, handlerFunc) - handler 函数签名必须是
func(*message.Message) ([]*message.Message, error),返回nil, nil表示仅消费,返回[]*message.Message, nil才触发下游发布
Kafka 场景下消息卡住或重复,先查 ConsumerGroup 和 Ack 配置
Watermill 把“何时算消费成功”完全交给底层驱动和你的 handler 返回值。Kafka 卡住、重复、反复重试,90% 是 ConsumerGroup 或 Ack 配置错位。
立即学习“go语言免费学习笔记(深入)”;
典型症状与解法:
-
context deadline exceeded错误 + 消息反复重试 → handler 执行超时触发 Kafka rebalance;解决:用ctx, cancel := context.WithTimeout(msg.Context(), 5*time.Second)包裹业务逻辑,并确保该 timeout 小于 Kafka 的session.timeout.ms(通常设为 10s 以内) - 所有实例都收到同一条消息 → 没设
ConsumerGroup,Kafka 认为每个 subscriber 是独立消费者;必须显式配置consumerGroup: "order-service"(Kafka)或queue: "order-queue"(AMQP) - 消息处理完却持续重试 → handler 返回了
error或 panic;记住:msg.Ack()只在 handler 返回nil时自动调用;若启用ManualAck模式,必须手动调msg.Ack(),否则 offset 永远不提交
防丢消息不能靠 Publisher,必须结合 Outbox 模式
kafka.NewPublisher 只是发送客户端,和数据库事务完全无关。订单创建后 DB commit 成功但 Kafka 网络失败,事件就丢了——这是生产环境最典型的“恰好一次”失效场景。
Outbox 模式是唯一可靠解法:
- 在业务数据库中建
outbox表,字段含id、topic、payload、published(bool) - DB 事务内:插入业务数据 + 插入 outbox 记录(
published = false) - 单独 goroutine 轮询
published = false的记录,调用Publisher.Publish()发送,成功后更新published = true - 轮询需带幂等性(比如用
SELECT ... FOR UPDATE SKIP LOCKED),避免多实例重复发
别试图用 Publisher 的重试配置代替 Outbox——网络抖动期间,重试可能发多次,而事务已提交,最终导致重复消费。
Google Cloud Pub/Sub 集成要配好 ProjectID 和订阅名生成规则
Watermill 对 Google Cloud Pub/Sub 的支持通过 github.com/ThreeDotsLabs/watermill-googlecloud/v2 提供,但配置稍有不同:Publisher 只需 ProjectID,Subscriber 必须明确指定订阅名生成逻辑。
关键点:
- Publisher 初始化:传
googlecloud.PublisherConfig{ProjectID: "my-project-id"}即可 - Subscriber 必须设置
GenerateSubscriptionName,否则默认用 topic 名作订阅名,多个服务实例会竞争同一订阅,导致消息被随机分发而非负载均衡;推荐写法:GenerateSubscriptionName: func(topic string) string { return "svc-order-" + topic } - 开发阶段可用
PUBSUB_EMULATOR_HOST环境变量对接本地模拟器,生产环境需设置GOOGLE_APPLICATION_CREDENTIALS指向服务账号密钥文件
容易被忽略的是:Google Cloud Pub/Sub 的订阅是全局资源,命名冲突会导致后续部署失败,GenerateSubscriptionName 必须包含服务标识和环境前缀(如 "prod-order-created"),不能只拼 topic。


















