必须用outbox.newpublisher绑定事务才能保证sql持久化事件流,因其将事件写入outbox表与业务数据同处一个事务;而kafka.newpublisher等不感知事务,无法保障一致性。

直接用 kafka.NewPublisher 或 gochannel.NewPublisher 发消息,和数据库事务是割裂的——订单写进 PostgreSQL,消息却因网络抖动丢了。要真正支持 SQL 持久化事件流,必须把消息发布行为绑进事务里,靠 Watermill 的 OutboxPublisher + 数据库驱动(如 watermill-sql)组合实现。
为什么不能只靠 Publisher 配置?
Watermill 的 kafka.NewPublisher、amqp.NewPublisher 都只是“发消息客户端”,不感知数据库事务。哪怕你先 tx.Commit() 再调 publisher.Publish(),中间出错(比如 Kafka 不可达、序列化失败),事件就彻底丢失,且无法回滚已提交的 DB 变更。
真正可靠的路径是:把事件先写进数据库的 outbox 表,再由独立协程轮询并投递——这样写 DB 和写 outbox 在同一事务中,原子性有保障。
-
gochannel.NewPublisher仅适合本地调试,无持久化能力 -
sql.NewPublisher(来自github.com/ThreeDotsLabs/watermill-sql)才是真正支持 SQL 持久化的发布器,但它本身也不自动绑定事务;必须配合outbox.NewPublisher封装使用 - 别手动往
outbox表 INSERT ——用outbox.NewPublisher,它会自动处理 schema、事务绑定、消息序列化
如何初始化带事务绑定的 OutboxPublisher?
关键不是选哪个底层 Pub/Sub,而是让 outbox.NewPublisher 拿到你的 *sql.Tx 或 *gorm.DB 实例。以 PostgreSQL + GORM 为例:
// 假设你已有 *gorm.DB 实例 db
outboxPub := outbox.NewPublisher(
db, // 注意:传的是 *gorm.DB,不是 *sql.DB
sql.NewPublisher(
sql.PublisherConfig{
Topic: "order_events",
},
logger,
),
logger,
)
// 在业务事务中调用
err := db.Transaction(func(tx *gorm.DB) error {
// 1. 写业务数据(如创建订单)
if err := tx.Create(&Order{ID: "o-123", Total: 100}).Error; err != nil {
return err
}
// 2. 发布事件(这行必须在 tx.Commit() 前!)
msg := message.NewMessage(uuid.NewString(), []byte(`{"order_id":"o-123"}`))
if err := outboxPub.Publish(ctx, "order_topic", msg); err != nil {
return err
}
return nil // tx.Commit() 由 gorm 自动触发
})
-
outbox.NewPublisher第一个参数必须是能执行Exec/Query的 DB 实例(*gorm.DB或*sql.DB),它会自动建outbox表(首次运行时) - 必须在事务 commit 前调
outboxPub.Publish(),否则事件不会落库 -
sql.NewPublisher是实际投递层,负责把 outbox 表里的记录发给 Kafka/RabbitMQ;它不参与事务,只读取已 commit 的 outbox 记录
Subscriber 怎么消费 SQL outbox 里的消息?
消费者端不需要改逻辑,但必须用 sql.NewSubscriber 替代 kafka.NewSubscriber 或 gochannel.NewSubscriber,否则无法从 outbox 表拉取消息:
subscriber, err := sql.NewSubscriber(
sql.SubscriberConfig{
Topic: "order_topic",
// 必须指定表名,且该表需由 outbox publisher 创建或提前建好
TableName: "outbox",
// 如果用 GORM,这里传 *gorm.DB;如果用 database/sql,传 *sql.DB
DB: db,
},
logger,
)
-
sql.NewSubscriber会轮询outbox表(默认每 100ms 查一次),把未处理的消息推给 Router - 它不依赖 Kafka/RabbitMQ 连接,纯 SQL 驱动,适合离线场景或强一致性要求的系统
- 注意:SQL Subscriber 不支持 ConsumerGroup 语义,多实例并发消费需靠数据库行锁或应用层去重(例如用
SELECT ... FOR UPDATE SKIP LOCKED)
容易被忽略的三个硬伤
很多人跑通了 Publish,但消费端卡住、重复、漏消息,问题往往不在 Watermill 本身,而在 SQL 层配置细节:
- 没建
outbox表或字段类型不匹配:outbox.NewPublisher默认建表,但若你手动建了,必须确保topic(VARCHAR)、payload(TEXT/BLOB)、created_at(TIMESTAMP)字段存在且类型兼容,否则Publish会静默失败 - Subscriber 启动后不消费:检查
sql.NewSubscriber的DB实例是否开启了 auto-commit(GORM 默认关,*sql.DB默认开),必须确保轮询查询不受事务隔离级别阻塞 - 消息重复投递:SQL Subscriber 默认不自动标记已处理,靠
DELETE FROM outbox WHERE id = ?完成“Ack”。若 handler panic 或进程崩溃,该 DELETE 未执行,下次轮询还会捞到同一条——这不是 bug,是设计使然;你要在 handler 里确保幂等,或启用sql.WithAckStrategy(sql.AckStrategyDelete)
最常被跳过的动作是:没确认 outbox 表是否真被创建、没验证 sql.NewSubscriber 轮询日志是否输出、没在 handler 里加 recover ——这些细节不显眼,但一出问题就卡死整条事件链。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











