使用 Watermill 与 Cloud Firestore 构建事务性 Pub/Sub:原理、配置与实战
【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill
导读
本文基于 Watermill 仓库官方文档 Firestore Pub/Sub 展开,系统讲解如何将 Google Cloud Firestore(一款云托管的 NoSQL 文档数据库)用作 Watermill 的消息总线。核心亮点在于 Firestore 提供普通Publisher与TransactionalPublisher两种发布器:后者允许你在同一个数据库事务里既写入业务数据又发布消息,从而从根本上规避"数据已保存但事件未发出"或"事件已发出但数据未保存"的一致性问题。读完本文,你将掌握 Firestore Pub/Sub 的安装方式、配置项、订阅命名机制与 Marshaler 工作原理,并理解如何与 Watermill 的 Forwarder 组件配合,实现经典的 Outbox 模式。
Cloud Firestore:为什么用它来做 Pub/Sub
Cloud Firestore 是 Google 提供的云托管 NoSQL 数据库,面向移动端、Web 与服务器应用,提供实时同步、离线支持与可扩展的文档模型。在绝大多数场景里,事件驱动系统会选择 Kafka、NATS、Google Pub/Sub 这类专用消息中间件,但 Firestore 有一个其他中间件难以替代的能力:支持在数据库事务中发布消息。
Watermill 官方对 Firestore Pub/Sub 的定位非常明确(见 firestore.md):
This Pub/Sub comes with two publishers. To publish messages in a transaction use the
TransactionalPublisher. If you do not want to publish messages in transaction use the normalPublisher.
也就是说,该 Pub/Sub 适配器内置了两个发布器:
TransactionalPublisher:在 Firestore 事务内发布消息,可与业务数据写入共用同一事务;- 普通
Publisher:不依赖事务的常规发布器。
选择 Firestore 而非专用消息系统,最典型的收益是一致性:当你需要在保存业务数据的同时发出消息时,如果两者不在同一事务内,就会陷入两难——消息发出但数据没保存,或者数据保存了但消息没发出。而借助TransactionalPublisher,数据与消息会被原子地一致持久化;事务提交之后,再通过订阅者把消息转发(relay)到其他更专用的 Pub/Sub 系统,形成"Firestore 落库 + 中间件分发"的分层架构。
特性一览
Watermill 官方文档以表格形式给出该 Pub/Sub 的特性支持情况,这是选择消息基础设施时最关键的评估依据:
| Feature | Implements | Note |
|---|---|---|
| ConsumerGroups | yes | |
| ExactlyOnceDelivery | no | |
| GuaranteedOrder | no | |
| Persistent | yes |
解读这份特性表:
- ConsumerGroups(消费者组):支持。可以像 Kafka 消费者组那样,让多个订阅者协同消费同一主题下的消息;
- ExactlyOnceDelivery(精确一次投递):不支持。这意味着在极端失败场景下,消息可能被重复投递,业务侧应具备幂等处理能力(可借助 Watermill 的 deduplicator 中间件 等方案兜底);
- GuaranteedOrder(消息有序):不支持。消息之间没有严格的全局顺序保证,对顺序敏感的业务需要在应用层自行处理(如使用 CQRS 组件中的有序事件方案);
- Persistent(持久化):支持。消息被写入 Firestore 文档,天然具备持久性,不会因进程崩溃而丢失。
安装
Firestore Pub/Sub 以独立的 Go 模块发布,安装命令与安装任意 Watermill 适配器一致:
go get github.com/ThreeDotsLabs/watermill-firestore说明:该模块属于 Watermill 生态的独立仓库,不在当前仓库源码目录内;当前仓库的 README.md 中也将其列为官方推荐的 Pub/Sub 适配器之一(见 "Firestore Pub/Sub (
github.com/ThreeDotsLabs/watermill-firestore)" 条目)。
先理解 Watermill 的 Pub/Sub 抽象
在深入 Firestore 适配器之前,有必要先回顾 Watermill 核心层定义的发布/订阅契约,因为这决定了Publisher、Subscriber、TransactionalPublisher各自承担的职责。契约定义在 message/pubsub.go:
type Publisher interface { // Publish publishes provided messages to the given topic. Publish(topic string, messages ...*Message) error // Close should flush unsent messages if publisher is async. Close() error } type Subscriber interface { // Subscribe returns an output channel with messages from the provided topic. Subscribe(ctx context.Context, topic string) (<-chan *Message, error) // Close closes all subscriptions with their output channels and flushes offsets etc. when needed. Close() error }从源码注释中可以提炼出对理解 Firestore 适配器至关重要的几点:
Publish可以是同步或异步的,取决于具体实现——Firestore 适配器内部就是对 Firestore 文档写入的封装;- 大多数 Publisher 实现不支持消息的原子发布,即一批消息中某一条失败时,后续消息可能不会被发布——这正是
TransactionalPublisher存在的意义; Subscribe返回消息通道,必须调用Ack()确认消费、Nack()触发重投——Firestore 订阅者同样遵循这套确认语义,以保证 at-least-once 风格的可靠消费。
Firestore Pub/Sub 之所以能在事务中发布消息,本质上是让Publisher接口的实现依赖一个外部的 Firestore 事务句柄(而非仅仅独立的客户端),从而与业务的数据写入共享同一个事务作用域。
配置:Publisher 与 Subscriber
Firestore Pub/Sub 的配置围绕发布端与订阅端分别展开,官方文档通过load-snippet-partial短代码,直接从外部仓库watermill-firestore的源码文件中截取配置结构体(src-link/watermill-firestore/pkg/firestore/publisher.go与subscriber.go),以保证文档与源码同步。
Publisher 配置
发布端配置结构体PublisherConfig struct定义于外部仓库pkg/firestore/publisher.go中。结合 Watermill 各数据库适配器的共性设计,可以推断其核心职责包括:
- 指定 Firestore 客户端、项目/数据库信息;
- 指定消息写入的集合(collection)路径;
- 指定事务获取方式(供
TransactionalPublisher使用); - 指定 Marshaler,控制 Watermill 消息如何转换为可存储的 Firestore 文档结构。
(注:具体字段名与默认值以外部仓库watermill-firestore的实际源码为准,本仓库不包含该模块源码,故不做逐字段断言。)
Subscriber 配置
订阅端配置结构体SubscriberConfig struct定义于外部仓库pkg/firestore/subscriber.go中。从官方文档的说明来看,其中有一个核心函数字段值得特别注意:
GenerateSubscriptionName func(topic string) string:用于生成订阅名称,默认实现是"主题名 +_sub后缀"。
其余字段(如轮询间隔、订阅处理并发度、Ack/Nack 策略等)同样以外部仓库源码为准。订阅端的核心机制如下:订阅者通过周期性查询 Firestore 中的消息文档来发现新消息,处理完成后写入 Ack 状态,从而在不依赖专用 broker 的情况下实现可靠消费。
订阅名称:一个决定消息如何分配的关键配置
官方文档用整整一节强调订阅名称的语义,这对正确使用 Firestore Pub/Sub 至关重要:
- 必须显式订阅才能收到消息。要向某主题接收消息,必须先为该主题创建订阅;只有订阅创建之后发布到该主题的消息才会被订阅者收到。这与 Kafka 等"保留全部消息、按 offset 消费"的模型有本质区别,更接近 Google Pub/Sub 的"订阅即消费起点"语义。
- 一个主题可以拥有多个订阅,但一个订阅只能归属一个主题。
- 订阅是自动创建的:在 Watermill 中,调用
Subscribe()时订阅会自动建立,无需手动操作。订阅名称由传给SubscriberConfig.GenerateSubscriptionName的函数生成,默认就是主题名加_sub后缀。
特别值得注意的一个实战场景(原文明确提示):如果你想让多个订阅者以不同的方式处理同一主题的消息,就必须使用自定义的GenerateSubscriptionName函数为每个订阅者生成唯一的订阅名称。因为默认命名规则下,多个订阅者订阅同一主题会共享同一个订阅名,从而只消费到部分消息(类似消费者组内分区分配);只有各自使用独立订阅名,每个订阅者才能各自收到该主题的完整消息流。
Marshaler:Watermill 消息与 Firestore 文档的桥梁
Watermill 的message.Message无法直接被 Firestore 存储——Firestore 只认它自己的文档结构。因此适配器引入Marshaler负责两者之间的转换:
Watermill's messages cannot be stored directly in Firestore. The marshaler is responsible for converting them to a type which can be stored by Firestore.
其接口定义同样来自外部仓库pkg/firestore/marshaler.go,职责涵盖:
- 将 Watermill 消息(
UUID、Payload、Metadata)转换为可写入 Firestore 的文档类型; - 反向将 Firestore 文档还原为 Watermill 消息(供订阅端消费)。
官方文档明确给出一个务实建议:默认实现足以覆盖绝大多数应用场景,你大概率不需要实现自定义 Marshaler。因此在实际项目中,直接使用默认 Marshaler 即可;只有当你的消息结构有特殊定制需求(例如需要额外存储自定义字段、或对文档结构有严格约束)时,才考虑基于该接口自行实现。
与 Forwarder 组合:落地 Outbox 模式
Firestore 事务性发布的最大价值,在于它可以成为Outbox(发件箱)模式的存储后端,这一点在仓库的 Forwarder 组件文档 中得到了明确佐证。该文档描述了一个经典的"抽奖"业务案例,并指出:
In case the database you're using is one among MySQL, PostgreSQL (or any other SQL), Firestore or Bolt, you can publish messages to them.Forwardercomponent will help you with picking all the messages you publish to the database and forwarding them to a message broker of yours.
其工作链路是:
- 业务代码在 Firestore 事务内同时写入业务数据,并通过
TransactionalPublisher发布事件消息——两者原子提交,从根上消除"数据与事件不一致"; - Forwarder 组件(见 components/forwarder/forwarder.go)作为后台守护进程,通过 Firestore 订阅者监听该中间主题上的消息;
- Forwarder 将包裹(envelope)解开后,把原始消息转发到真正面向外部的消息 broker(如 Google Pub/Sub、Kafka、NATS 等),供下游服务消费。
这正好呼应了本文开头引用的官方表述:"After transactionally publishing messages in Firestore you can then subscribe to them and relay them to a different Pub/Sub system"(事务性发布消息到 Firestore 后,你可以订阅这些消息并把它们转发到另一个 Pub/Sub 系统)。
使用要点与边界
综合上述内容,使用 Watermill + Firestore Pub/Sub 时建议关注以下几点:
- 优先事务性发布:凡是"业务数据 + 事件消息"需要保持一致性的场景,一律使用
TransactionalPublisher,不要依赖"先发消息、后存数据"或"先存数据、后发消息"的顺序补救方案——前者会因事件先发出而数据未落库造成下游误判,后者会因消息丢失导致下游无感知; - 为每个不同消费语义配置独立订阅名:默认
_sub后缀只适合单一消费者组;多订阅者差异化处理同一主题时,务必自定义GenerateSubscriptionName; - 接受非精确一次投递与无序:特性表明确不支持 ExactlyOnceDelivery 与 GuaranteedOrder,业务侧需要做好幂等与乱序容忍设计(可参考 message/router/middleware/deduplicator.go 等现成中间件);
- 明确适用边界:如果你不需要"事务内发布"这个能力,直接用 Kafka、Google Pub/Sub 等专用中间件通常是更优选择;Firestore 方案的典型形态是"Firestore 作为事务性落点 + Forwarder 转发到专用 broker",而非替代专用 broker 成为全链路消息骨干。
总结
Watermill 的 Firestore Pub/Sub 适配器将 Google 云数据库与事件驱动架构优雅地结合在一起:凭借TransactionalPublisher的事务内发布能力,业务数据与领域事件可以在同一事务中原子持久化,彻底规避顺序发布带来的数据一致性陷阱;配合订阅名称机制、默认 Marshaler 以及 Forwarder 组件的 Outbox 转发链路,即可构建一套"本地一致、对外可靠分发"的生产级事件驱动方案。本文对应的原始官方文档位于 docs/content/pubsubs/firestore.md,如需查阅所有受支持的 Pub/Sub 适配器清单,可参考 pubsubs 文档索引。
【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考