ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

使用 Watermill 与 Cloud Firestore 构建事务性 Pub/Sub:原理、配置与实战

使用 Watermill 与 Cloud Firestore 构建事务性 Pub/Sub:原理、配置与实战 使用 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/SubCloud Firestore 是 Google 提供的云托管 NoSQL 数据库面向移动端、Web 与服务器应用提供实时同步、离线支持与可扩展的文档模型。在绝大多数场景里事件驱动系统会选择 Kafka、NATS、Google Pub/Sub 这类专用消息中间件但 Firestore 有一个其他中间件难以替代的能力支持在数据库事务中发布消息。Watermill 官方对 Firestore Pub/Sub 的定位非常明确见 firestore.mdThis Pub/Sub comes with two publishers. To publish messages in a transaction use theTransactionalPublisher. 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 的特性支持情况这是选择消息基础设施时最关键的评估依据FeatureImplementsNoteConsumerGroupsyesExactlyOnceDeliverynoGuaranteedOrdernoPersistentyes解读这份特性表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.gotype 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 与 SubscriberFirestore 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函数为每个订阅者生成唯一的订阅名称。因为默认命名规则下多个订阅者订阅同一主题会共享同一个订阅名从而只消费到部分消息类似消费者组内分区分配只有各自使用独立订阅名每个订阅者才能各自收到该主题的完整消息流。MarshalerWatermill 消息与 Firestore 文档的桥梁Watermill 的message.Message无法直接被 Firestore 存储——Firestore 只认它自己的文档结构。因此适配器引入Marshaler负责两者之间的转换Watermills 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 youre 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),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进