尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Encore Go 事务性 Pub/Sub Outbox 实战:用 Outbox 表与 Relay 保证数据库与订阅者之间的最终一致性

Encore Go 事务性 Pub/Sub Outbox 实战:用 Outbox 表与 Relay 保证数据库与订阅者之间的最终一致性 Encore Go 事务性 Pub/Sub Outbox 实战用 Outbox 表与 Relay 保证数据库与订阅者之间的最终一致性【免费下载链接】encoreThe infrastructure platform for the intelligence era项目地址: https://gitcode.com/GitHub_Trending/encor/encore构建事件驱动应用时最难解决的问题之一就是服务之间的一致性每个服务通常拥有自己的数据库并通过 Pub/Sub 把业务事件通知给其他系统而消息发布与数据库写入天然不处于同一个事务中一旦二者之间出现失败就会产生数据不一致。本文基于 Encore 官方指南 docs/go/how-to/pubsub-outbox.md完整讲解事务性 OutboxTransactional Outbox模式的落地方式如何把 Pub/Sub Topic 绑定到数据库事务、如何用x.encore.dev/infra/pubsub/outbox包完成消息的暂存与转发以及如何通过 Relay 把已提交的消息发布到真实 Topic。读完本文你将能够在不改变既有发布代码的前提下为 Encore Go 服务接入一套“数据库写 事件发布”原子一致、且支持可插拔存储后端的事务性 Outbox 方案。一、问题背景为什么数据库写入与 Pub/Sub 发布会不一致在一个典型的微服务架构中每个服务都拥有自己独立的数据库同时通过 Pub/Sub 向其他系统广播业务事件如用户注册、订单创建等。最常见的写法是在数据库事务中写入业务数据在事务提交后或提交前调用topic.Publish发布事件。这种“先写库、再发消息”或“先发消息、再写库”的方式都会引入不确定性如果先发布消息、后提交数据库事务订阅者可能在数据尚未落库时就消费到事件如果先提交事务、后发布消息则消息发送失败时事件会永久丢失事务回滚与消息发布之间没有原子性保证两者总会出现部分成功的情况。这正是 Outbox 模式要解决的问题将“发布消息”这一动作本身转化为一次数据库写入与业务数据在同一事务内提交从而让“写业务数据”和“写事件”具备原子性。二、Encore 的事务性 Outbox核心工作方式Encore 在x.encore.dev/infra/pubsub/outbox包中提供了对事务性 Outbox 模式的官方支持。它的核心思路可以概括为两句话写入侧Outbox 表把一个 Pub/Sub Topic 绑定到一个数据库事务上之后对该 Topic 的所有topic.Publish调用都会被翻译成向outbox表中插入一行记录而不是真正发布消息转发侧Relay当事务提交后消息可以被一个Relay拾取——Relay 持续轮询outbox表把新增的行逐条发布到真正的 Pub/Sub Topic 上订阅者照常消费。用一张时序来描述整个生命周期业务代码调用 ref.Publish(msg) │ ▼ 插入 outbox 表一行topic、data、inserted_at │ ▼ 数据库事务 COMMIT业务数据 outbox 行原子生效 │ ▼ Relay 轮询 outbox 表发现新行 │ ▼ Relay 调用真实 topic.Publish(msg) │ ▼ 订阅者正常收到事件只要事务没有提交outbox表中的行就不可见、不会被 Relay 拾取一旦提交消息必然会在后续的轮询中被转发从而同时保证了“不丢失”与“不提前可见”。注意常规非 Outbox用法下topic.Publish返回的消息 id 与订阅者处理消息时收到的 id 是同一个。而使用 Outbox 时由于真实消息 id 要到事务提交、消息真正发布时才会产生因此topic.Publish返回的 id 引用的是outbox 表中的行 id。如果代码依赖该返回值来关联消息需要留意这一语义差异。三、发布消息到 OutboxTopic 引用 事务绑定3.1 先理解 Topic 引用TopicRef要把 Topic 绑定到 Outbox文档推荐的方式是使用Pub/Sub Topic 引用topic references。为什么要用 TopicRef 而不是直接操作*pubsub.Topic关键在于 Encore 的静态分析机制。Encore 通过静态分析来确定“哪个服务向哪个 Topic 发布了消息”并据此完成基础设施编排、架构图渲染和 IAM 权限配置详见 docs/go/primitives/pubsub.md#using-topic-references。这意味着*pubsub.Topic变量不能随意在代码中传递否则静态分析将无法追踪。而pubsub.TopicRef允许你获得一个可以自由传递、且按需声明权限的引用signupRef : pubsub.TopicRef[pubsub.Publisher[*SignupEvent]](Signups) // signupRef 的类型是 pubsub.Publisher[*SignupEvent]只允许发布操作从本仓库的运行时源码可以印证这一点runtimes/go/pubsub/refs.go 中定义了权限接口type Publisher[T any] interface { // Publish publishes a message to the topic. Publish(ctx context.Context, msg T) (id string, err error) // Meta returns metadata about the topic. Meta() TopicMeta }而TopicRef的实现runtimes/go/pubsub/refs.go#L31-L49要求类型参数P必须是一个“声明权限的接口”即实现TopicPerms[T]的接口如pubsub.Publisher[T]返回的引用被收窄到该权限范围内func TopicRef[P TopicPerms[T], T any](topic *Topic[T]) P { return any(topicRef[T]{Topic: topic}).(P) }这正是 Outbox 绑定得以成立的基础outbox.Bind接收的正是这种pubsub.Publisher[...]类型的引用绑定后你仍然拥有与普通 Topic 完全一致的接口和类型安全已有发布代码无需任何改动。3.2 将 Topic 绑定到事务outbox.Bind TxPersister绑定动作本身一行代码即可完成。完整示例用户创建后通知订阅者-- outbox.go -- // 创建一个 SignupsTopic略去具体配置。 var SignupsTopic pubsub.NewTopic*SignupEvent // 创建带发布权限的 topic ref。 ref : pubsub.TopicRef[pubsub.Publisher[*SignupEvent]](SignupsTopic) // 将 ref 绑定到事务性 outbox。 import x.encore.dev/infra/pubsub/outbox var tx *sqldb.Tx // 从业务代码中获取一个数据库事务 ref outbox.Bind(ref, outbox.TxPersister(tx)) // 此后所有 ref.Publish() 调用都会向 outbox 表插入一行。核心 API 说明outbox.Bind(ref, persister)接收一个pubsub.Publisher[T]类型的 Topic 引用和一个“持久化器”persister返回一个绑定后的新引用outbox.TxPersister(tx)把当前数据库事务包装成持久化器。绑定后Publish不再真正发布而是在该事务内插入 outbox 行类型安全保持绑定前后的引用都实现了同一个pubsub.Publisher[T]接口因此调用方如业务逻辑感知不到任何差异。3.3 可插拔的存储后端outbox.Bind的持久化器是可插拔的这意味着 Outbox 模式可以配合任意“支持事务的存储后端”使用。官方开箱即用地提供了以下实现Encore 自带的encore.dev/storage/sqldb包即outbox.TxPersister配合sqldb.Tx标准库database/sql驱动github.com/jackc/pgx/v5驱动。对于其他数据库你可以基于PersistFunc编写自己的持久化器详见该包的 Go 参考文档。这意味着无论你用的是 Encore 管理的数据库还是通过database/sql/pgx/v5接入的已有 PostgreSQL都能使用同一套 Outbox 机制。3.4 数据库迁移outbox 表结构使用 SQL 存储后端时数据库中必须存在对应的outbox表。官方给定的建表迁移脚本如下-- db_migration.sql -- -- The database used must contain the below database table: CREATE TABLE outbox ( id BIGSERIAL PRIMARY KEY, topic TEXT NOT NULL, data JSONB NOT NULL, inserted_at TIMESTAMPTZ NOT NULL ); CREATE INDEX outbox_topic_idx ON outbox (topic, id);各列职责列名类型说明idBIGSERIAL主键自增topic.Publish返回的“消息 id”实际上就是这一行记录的 idtopicTEXT目标 Topic 的名称Relay 靠它区分消息应该发布到哪个 TopicdataJSONB消息体SignupEvent等负载的序列化结果使用 JSONB 便于存储与查询inserted_atTIMESTAMPTZ插入时间Relay 可据此进行增量轮询索引outbox_topic_idx建在(topic, id)上与“按 Topic 过滤 按 id 递增排序”的轮询模式相匹配。当事务提交后通过绑定引用ref发布的所有消息都会以行记录的形式落在outbox表中等待被 Relay 消费。四、从 Outbox 消费消息Relay 轮询与发布4.1 Relay 的工作原理事务提交后消息已经安全地躺在outbox表中接下来需要把它们发布到真正的 Pub/Sub Topic。这一步由outbox.Relay完成Relay 持续轮询outbox表将任何新增的消息发布到对应的真实 Topic。Relay 同样支持可插拔的存储后端官方默认提供的是基于 Encore 内置 SQL 数据库支持的outbox.SQLDBStore(db)你也可以为其编写其他数据库的实现。4.2 注册 Topic 并启动轮询需要被轮询的 Topic 必须显式注册到 Relay 上这一般在服务初始化阶段完成。完整的服务端示例-- user/service.go -- package user import ( context encore.dev/pubsub encore.dev/storage/sqldb x.encore.dev/infra/pubsub/outbox ) type Service struct { signupsRef pubsub.Publisher[*SignupEvent] } // db 是存放 outbox 表的数据库。 var db sqldb.NewDatabase(...) // 创建 SignupsTopic略去具体配置。 var SignupsTopic pubsub.NewTopic*SignupEvent func initService() (*Service, error) { // 初始化 Relay从我们的数据库轮询。 relay : outbox.NewRelay(outbox.SQLDBStore(db)) // 注册 SignupsTopic使其进入被轮询的列表。 signupsRef : pubsub.TopicRef[pubsub.Publisher[*SignupEvent]](SignupsTopic) outbox.RegisterTopic(relay, signupsRef) // 启动轮询。 go relay.PollForMessage(context.Background(), -1) return Service{signupsRef: signupsRef}, nil }拆解这段初始化流程outbox.NewRelay(outbox.SQLDBStore(db))以指定的 SQL 数据库为后端创建 Relayoutbox.RegisterTopic(relay, signupsRef)把“要轮询的 Topic”注册进 RelayRelay 只处理注册过的 Topic 对应的 outbox 行go relay.PollForMessage(context.Background(), -1)以 goroutine 形式启动持续轮询第二个参数用于控制轮询节奏示例传入-1表示无额外休眠、尽可能快地连续轮询context.Background()作为轮询循环的生命周期上下文。启动之后每次ref.Publish产生的 outbox 行都会在事务提交后的下一次轮询中被 Relay 拾取并发布到真实的SignupsTopic所有注册的订阅者随即收到事件。业务代码侧自始至终只面向pubsub.Publisher[*SignupEvent]接口编程与普通 Pub/Sub 完全一致。五、源码与仓库证据Outbox 能力在 Encore 生态中的定位虽然x.encore.dev/infra/pubsub/outbox是一个独立发布的扩展包但它在 Encore 生态中是被官方文档与运行时共同支撑的一等公民能力在本仓库中可以找到多处理论与实现层面的证据Topic 引用机制是绑定的基础runtimes/go/pubsub/refs.go 完整定义了TopicPerms、Publisher与TopicRef正是outbox.Bind与outbox.RegisterTopic所操作的引用类型也解释了为什么绑定后的代码可以绕过 Encore 对*pubsub.Topic的静态分析限制、被自由传递到任意库代码或 service structs 中Outbox 模式被官方文档反复引用除了本篇 pubsub-outbox.mddocs/go/primitives/pubsub.md#using-topic-references 对 TopicRef 的权限声明与基础设施编排语义做了前置讲解docs/go/develop/testing.md#L105-L106 还专门提到该包“定义了一个仅用于测试的数据库用于对 outbox 功能做集成测试”说明官方为 Outbox 的集成测试场景提供了开箱支持命令行工具链中的编排元数据go_llm_instructions.txt中同样收录了outbox.Bind、TxPersister、CREATE TABLE outbox、NewRelay/RegisterTopic等完整示例go_llm_instructions.txt#L718-L743可作为快速检索的索引也佐证了这一模式属于官方推荐的 Go 应用开发范式。六、实践要点与适用边界在把事务性 Outbox 引入实际项目时值得记住以下几点绑定时机outbox.Bind接收的tx *sqldb.Tx必须来自当前正在执行的业务事务——绑定的语义是“本次事务内的所有 Publish 都写入 outbox”因此绑定应与事务生命周期对齐通常在开启事务后、执行业务写入前完成。消息可见性事务未提交前outbox 行对 Relay 不可见因此订阅者绝不会提前消费到“尚未提交”的事件事务回滚时outbox 行随之一并回滚事件自然不会被发布——这正是该模式一致性保证的来源。消息 id 语义变化Outbox 模式下Publish返回的是 outbox 行 id 而非真实消息 id若现有代码把返回值当作真实消息 id 使用如写入日志、关联追踪需要相应调整。轮询延迟Relay 采用轮询方式转发消息端到端延迟取决于轮询间隔适合对实时性要求不苛刻的业务事件如需接近实时的体验可调小轮询间隔。存储后端选择优先使用官方开箱即用的sqldb/database/sql/pgx/v5实现其他存储可基于PersistFunc扩展但需要自行保证持久化行为与事务语义的正确性。综上Encore 的事务性 Outbox 以“绑定 Topic 引用 事务内写 outbox 表 Relay 轮询转发”三步用极小的侵入成本化解了事件驱动架构中最棘手的“库表与消息不一致”问题业务代码保持原有发布接口不变一致性由数据库事务与 Relay 共同兜底。无论是新项目还是需要为既有 Encore Go 服务补齐一致性保证的场景这都是值得优先考虑的官方方案。【免费下载链接】encoreThe infrastructure platform for the intelligence era项目地址: https://gitcode.com/GitHub_Trending/encor/encore创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表