news 2026/9/15 12:44:22

Encore Go 后端 Pub/Sub 完全指南:用声明式 Topic 与 Subscription 构建异步事件驱动系统

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Encore Go 后端 Pub/Sub 完全指南:用声明式 Topic 与 Subscription 构建异步事件驱动系统

Encore Go 后端 Pub/Sub 完全指南:用声明式 Topic 与 Subscription 构建异步事件驱动系统

【免费下载链接】encoreThe infrastructure platform for the intelligence era项目地址: https://gitcode.com/GitHub_Trending/encor/encore

本文基于 Encore 开源仓库 docs/go/primitives/pubsub.md 及 runtimes/go/pubsub 运行时源码,系统讲解如何用 Encore Backend Framework 在 Go 应用中通过 Pub/Sub 以声明式、云无关(cloud-agnostic)的方式构建异步消息系统。你将掌握 Topic 的定义与三种投递语义(at-least-once / exactly-once / 有序投递)、事件的发布与 TopicRef 引用、订阅的配置与错误重试/死信队列机制,以及如何在单元测试中隔离并断言消息发布行为——这些能力可直接用于解耦服务、提升系统可靠性与响应速度。

为什么需要 Pub/Sub:从同步 API 调用到事件广播

Pub/Sub(发布/订阅)是一种让系统组件通过异步广播事件进行通信的架构模式。Encore 的 Backend Framework 让开发者以纯声明式的方式使用 Pub/Sub:在部署时,Encore 会自动为你配置所需的云基础设施(AWS、GCP、自有云等),本地开发则使用内置的 NSQ 模拟实现。开发者无需编写任何云厂商 SDK 代码。

API 直连方式的痛点

以用户注册为例:假设注册成功后需要发送欢迎邮件并在分析系统中记录注册信息。如果只用 API 调用,调用链如下:

  1. user服务开启数据库事务,写入用户记录;
  2. user服务调用email服务发送欢迎邮件;
  3. email服务再调用邮件供应商真正发送邮件;
  4. 邮件发送成功后,email服务回复user服务;
  5. user服务再调用analytics服务记录注册信息;
  6. analytics服务写入数据仓库;
  7. analytics服务回复user服务;
  8. user服务提交数据库事务;
  9. user服务才回复用户“注册成功”。

这种设计有两个明显问题:响应时间被下游拖累——如果邮件供应商耗时 3 秒,用户就要等 3 秒才能收到注册成功响应,而实际上用户写入数据库那一刻就可以确认注册;故障爆炸半径大——如果数据仓库故障,所有用户注册都会失败,而分析系统是纯内部功能,不应该影响用户注册。

Pub/Sub 方式的优势

改用 Pub/Sub 后,注册流程变为:

  1. user服务开启数据库事务,写入用户记录;
  2. signupsTopic 发布一条注册事件;
  3. 提交事务,立即回复用户注册成功。

此时user服务与emailanalytics服务彻底解耦:后两者各自订阅signupsTopic 并行处理事件;若处理失败,事件会自动退避重试,直到成功或达到最大重试次数后进入死信队列(DLQ)。user服务甚至完全不知道emailanalytics的存在,未来新增任何关注“新用户注册”的系统都无需改动user服务。这正是 Pub/Sub 提升可靠性(缩小故障爆炸半径)、加快用户响应速度、并通过反转服务间依赖降低开发认知负担的核心价值。

创建 Topic:声明事件流的核心

Topic(主题)是 Pub/Sub 的核心,它是一个命名的、用于发布事件的通道。Topic 必须声明为包级变量,不能在函数内部创建。无论 Topic 定义在哪个服务中,任何服务都可以向它发布事件、任何服务也都可以订阅它。

创建 Topic 时需要指定:事件类型、唯一名称,以及定义其行为的配置。以“用户注册事件”为例:

package user import "encore.dev/pubsub" type SignupEvent struct{ UserID int } var Signups = pubsub.NewTopic*SignupEvent

TopicConfig 配置项

从 runtimes/go/pubsub/internal/types/public.go 的定义看,TopicConfig包含两个字段:

字段类型必填说明
DeliveryGuaranteeDeliveryGuarantee投递语义:AtLeastOnceExactlyOnce
OrderingAttributestring作为排序键的消息属性名,设置后同一键值的消息按发布顺序投递

运行时(topic.go)在创建 Topic 时会根据当前运行环境选择实现:单元测试场景使用internal/test的测试实现,未注册的场景降级为 noop 实现,正常部署时则按PubsubProviders匹配对应的云厂商实现(仓库内置 AWS、GCP、Azure、Encore Cloud 与本地 NSQ 五种 provider)。

投递语义:从 At-Least-Once 到 Exactly-Once

At-least-once(至少一次投递)

上面的示例配置保证:对于每个订阅,事件至少被投递一次。如果 Topic 认为事件未被成功处理,会尝试再次投递。

因此,所有订阅处理函数都应设计为幂等的——即处理函数被调用两次或多次时,从外部看与调用一次没有差别。实现幂等通常有两种方式:用数据库记录该事件触发的动作是否已执行;或者确保动作本身天然具备幂等性。

Exactly-once(精确一次投递)

DeliveryGuarantee设置为pubsub.ExactlyOnce,可在基础设施层面提供更强的保证,最小化消息被重复投递的可能性:

var Signups = pubsub.NewTopic*SignupEvent

但即便如此,仍有极少数情况下消息会被重投:例如网络问题导致“处理成功”的确认消息在到达云厂商之前丢失(即著名的两军问题 也明确指出:Exactly-once 只约束“投递给消费者”这一环节,不包含发布侧的去重——如果应用逻辑中Publish被调用了两次(例如应用层重试),消息会被投递两次;且 Exactly-once Topic 上的订阅相比 At-least-once 有更高的投递延迟。

启用 exactly-once 后,云厂商会施加吞吐限制:

  • AWS:Topic 每秒最多 300 条消息(参见 AWS SQS Quotas);
  • GCP:区域范围内所有 Topic 合计至少每秒 3,000 条消息(视区域可能更高,参见 GCP Pub/Sub Quotas)。

有序 Topic:保证同一实体的消息顺序

Topic 默认无序,消息可以按任意顺序投递,这允许并行处理以获得更高吞吐。但在某些场景下,针对某个特定实体,消息必须按发布顺序投递。

创建有序 Topic 的方式:将OrderingAttribute设置为事件类型某个顶层字段的pubsub-attr标签值。该字段值相同的消息,会按发布顺序投递给同一订阅者;排序键不同的消息之间顺序不受约束。

package example import ( "context" "encore.dev/pubsub" ) type CartEvent struct { ShoppingCartID int `pubsub-attr:"cart_id"` Event string } var CartEvents = pubsub.NewTopic*CartEvent func Example(ctx context.Context) error { // 这三条消息购物车 ID 相同,会按顺序投递 CartEvents.Publish(ctx, &CartEvent{ShoppingCartID: 1, Event: "item_added"}) CartEvents.Publish(ctx, &CartEvent{ShoppingCartID: 1, Event: "checkout_started"}) CartEvents.Publish(ctx, &CartEvent{ShoppingCartID: 1, Event: "checkout_completed"}) // 这条消息购物车 ID 不同,可能在任意时刻投递 CartEvents.Publish(ctx, &CartEvent{ShoppingCartID: 2, Event: "item_added"}) }

有序投递的注意点

  • 队头阻塞(head-of-line blocking):为维护顺序,同一排序键的消息必须等最早的消息处理完成或被送入死信队列后才继续投递,这可能在键值上堆积延迟。源码 internal/types/public.go 特别提醒:排序键下的消息处理出错时会先重试、再投递后续消息,因此配置重试策略时要充分考虑失败模式,避免积压;应通过健全的日志、告警和合适的订阅重试策略来缓解。
  • 本地环境无排序效果OrderingAttribute在本地开发环境中目前不生效。
  • 吞吐限制:各云厂商对有序 Topic 有吞吐限制——AWS为 Topic 每秒 300 条消息;GCP为每个排序键 1 MB/s(参见 GCP Pub/Sub Resource Limits)。

发布时排序键的提取发生在运行时 topic.go:Publish会先从消息中提取pubsub-attr标签标记的属性,若OrderingAttribute配置的属性不存在或为空字符串,会返回InvalidArgument错误(该情况理论上已被静态分析拦截,属于防御性检查)。

发布事件

发布事件只需在 Topic 上调用Publish,传入事件对象(即pubsub.NewTopic[Type]构造器指定的类型):

messageID, err := Signups.Publish(ctx, &SignupEvent{UserID: id}) if err != nil { return err } // 走到这里说明事件已成功发布, // 所有已注册的订阅者都会收到该事件。 // messageID 是消息的唯一 ID, // 订阅者处理事件时也会获得该 ID。

Signups声明为导出的包级变量后,其他服务也可以同样方式向该 Topic 发布事件。

深入Publish的运行时行为

从 topic.go 的实现看,一次Publish调用背后包含完整的处理链路:

  1. 上下文与合法性检查:context 已取消则直接返回;Topic 未通过NewTopic创建则返回Unimplemented错误;
  2. 属性提取与 JSON 序列化:通过utils.MarshalFields提取pubsub-attr标记的属性,并将消息体序列化为 JSON(序列化失败返回InvalidArgument);
  3. 排序键解析:若配置了OrderingAttribute,从属性中取出排序键;
  4. 追踪与关联 ID 传播:若当前处于某个请求中,会向属性注入encore_parent_trace_idencore_ext_correlation_id(外部关联 ID 优先,否则用 trace ID)以及平台请求的强制追踪标记,让订阅者可以把 trace 标记为当前请求的子 span——这正是 Encore 分布式追踪在异步消息链路上保持连贯的实现机制;
  5. 限流与发布:等待发布限流器(publishLimiter)放行后调用对应云厂商的PublishMessage真正发布;
  6. 错误归一化:发布失败统一包装为Unavailable错误码返回。

注意:Publish会阻塞直到消息被 Topic 成功接收;若返回错误,通常意味着发布失败,但不排除消息仍可能被订阅者收到(分布式系统的经典语义)。

使用 TopicRef:突破静态分析的引用限制

Encore 通过静态分析确定哪些服务向哪些 Topic 发布消息,据此正确配置基础设施、渲染架构图并配置 IAM 权限。因此*pubsub.Topic变量不能随意传递——那会使静态分析在许多场景下失效。

为此 Encore 提供了TopicRef,它返回一个可以自由传递的 Topic 引用:

signupRef := pubsub.TopicRef[pubsub.Publisher[*SignupEvent]](Signups) // signupRef 的类型是 pubsub.Publisher[*SignupEvent],仅允许发布操作

TopicRefTopic的关键区别是:引用必须预先声明所需权限,Encore 假定你声明的所有权限都会被使用。例如上面声明了pubsub.Publisher权限,Encore 就认为该服务会向 Topic 发布消息,并为此配置基础设施。从 refs.go 的实现看,TopicRef借助 Go 泛型接口(Publisher[T]接口约束为TopicPerms[T])把*Topic[T]收窄为只暴露声明权限的接口对象。

注意:TopicRef必须在服务内部声明,但引用本身可以自由传递给库代码、通过依赖注入注入到服务结构体中等。

订阅事件

创建订阅需调用pubsub.NewSubscription,同样以包级变量形式声明。每个订阅需要:

  • 要订阅的 Topic;
  • 一个在该 Topic 下唯一的名称;
  • 一个配置对象,其中至少包含处理事件的Handler函数。

订阅示例(放在email服务中):

package email import ( "encore.dev/pubsub" "user" ) var _ = pubsub.NewSubscription( user.Signups, "send-welcome-email", pubsub.SubscriptionConfig[*SignupEvent]{ Handler: SendWelcomeEmail, }, ) func SendWelcomeEmail(ctx context.Context, event *SignupEvent) error { // 发送邮件... return nil }

订阅可以定义在 Topic 所在的服务,也可以定义在应用的任何其他服务中。每个订阅都独立于同一 Topic 的其他订阅接收事件:如果某个订阅处理缓慢,它只会积压自己的未处理事件,其他订阅仍然实时处理新发布的事件。

订阅命名规范(见 subscription.go 的文档注释):名称必须在 Topic 内唯一,使用 kebab-case(小写字母数字与连字符),以字母开头、以字母或数字结尾,最长 63 个字符。部署后切勿更改订阅名或 Topic 名,否则在途消息可能丢失。

Handler 与 AckDeadline

传给 Handler 的ctx会在订阅的AckDeadline到达时被取消——这是消息被认为处理超时、可以被重新投递给其他订阅者的时间点。若不显式配置,默认值为30 秒(运行时 subscription.go 中AckDeadline == 0时默认设为 30 秒)。

SubscriptionConfig.Handler的文档(types.go)可以明确处理协议:Handler 应阻塞直到与消息相关的全部处理完成才返回;返回 nil 表示消息被确认(ack),不应再投递;返回非 nil 错误表示消极确认(nack),会触发重试(除非达到RetryPolicy.MaxRetries)。

基于服务结构体方法的 Handler

使用服务结构体做依赖注入时,通常希望把订阅处理函数定义为服务结构体的方法,以便访问注入的依赖。此时使用pubsub.MethodHandler

//encore:service type Service struct { /* ... */ } func (s *Service) SendWelcomeEmail(ctx context.Context, event *SignupEvent) error { // ... } var _ = pubsub.NewSubscription( user.Signups, "send-welcome-email", pubsub.SubscriptionConfig[*SignupEvent]{ Handler: pubsub.MethodHandler((*Service).SendWelcomeEmail), }, )

注意:pubsub.MethodHandler只允许引用服务结构体类型上的方法,不能是其他类型。其实现(subscription.go)本身是一个哨兵函数——真正的调用会在代码生成阶段被替换为初始化服务结构体的生成代码,因此该函数在运行时绝不会真正执行。

SubscriptionConfig 完整配置

创建订阅时,可通过SubscriptionConfig配置消息保留时长、重试策略等行为(types.go):

字段默认值说明
Handler必填处理消息的函数(或MethodHandler包装的服务方法)
MaxConcurrency视云厂商而定每实例同时处理的最大消息数;负数表示不限制;注意按实例计算(10 个实例 × 10 并发 = 100 并发);GCP Cloud Run 推送订阅与 Encore Cloud 环境不生效
AckDeadline30 秒消费者处理一条消息的最长时限,至少 1 秒
MessageRetention7 天未投递消息在 Topic 上保留多久后被清除
RetryPolicyMaxRetries: 100处理出错时的重试策略

重要约束SubscriptionConfig的所有字段必须是编译期常量,不能用函数调用表达式定义。这是 Encore 在部署时理解订阅确切需求、以便配置正确基础设施的前提。

RetryPolicy的定义在 internal/types/public.go:

字段默认值说明
MinBackoff10 秒两次重试之间的最小等待时间,不可为负
MaxBackoff10 分钟两次重试之间的最大等待时间,不可为负
MaxRetries1000 时使用默认值 100;>0 时表示重试 n 次后转入死信队列;pubsub.NoRetries(-2)表示任何错误/panic 立即进死信队列;pubsub.InfiniteRetries(-1)表示永远重试、不进入死信队列

各值会在编译期被解析以支撑云资源预配置,实际部署时可能被目标云厂商钳制到其支持的范围。

错误处理与死信队列(DLQ)

如果订阅处理函数返回错误,事件会按该订阅配置的重试策略进行重试。达到MaxRetries后,事件被放入该订阅者的死信队列(DLQ)。这样订阅可以继续处理后续事件,直到导致失败的 bug 被修复;修复后,DLQ 中的消息可以被手动释放、重新交给订阅者处理。

运行时还提供了兜底保护(subscription.go):如果 Handler 发生panic,会被包装为Internal错误码返回,从而触发重试/DLQ 流程而不是让整个进程崩溃。

测试 Pub/Sub:隔离、确定性与断言

Encore 使用特殊的测试实现来运行 Pub/Sub Topic(实现位于 runtimes/go/pubsub/internal/test/topic.go)。运行测试时,Topic 知道当前是哪个测试在运行,从而提供以下保证:

  • 订阅不会被触发:测试中发布的事件不会触发订阅者,使你能够独立测试发布者的行为,而不受订阅者副作用干扰;
  • 消息 ID 确定性:发布时生成的 Message ID 是确定性的(基于发布顺序),因此你的断言可以直接利用这一点;
  • 测试间隔离:每个测试与其他测试相互隔离,即使使用并行测试,一个测试发布的事件也不会影响其他测试。

Encore 提供辅助函数et.Topic访问测试 Topic,通过PublishedMessages()提取测试期间发布的事件:

package user import ( "testing" "encore.dev/et" "github.com/stretchr/testify/assert" ) func Test_Register(t *testing.T) { t.Parallel() ... 调用 Register() 并断言数据库变更 ... // 获取本测试中发布到 Signups Topic 的所有消息 msgs := et.Topic(Signups).PublishedMessages() assert.Len(t, msgs, 1) }

测试 Topic 还支持按需启用订阅者(TestTopic.PublishMessage会先记录消息,若测试实例启用了订阅,则异步触发对应订阅回调,模拟真实系统的发布行为),默认行为是订阅者关闭。

服务间一致性:事务性 Outbox 模式

事件驱动应用中保证服务间一致性颇具挑战,尤其是当数据库写入与 Pub/Sub 发布不在同一事务中时,可能导致服务间数据不一致。在不引入过多复杂度的前提下,推荐采用事务性 Outbox(transactional outbox)模式来解决。Encore 提供了完整的实现指南,参见 Pub/Sub Outbox 指南。

总结

Encore 的 Pub/Sub 抽象把“声明式 API”与“云无关的底层实现”结合起来:开发者只需用pubsub.NewTopic声明 Topic、用pubsub.NewSubscription声明订阅,Encore 的静态分析便会自动完成云基础设施的配置、IAM 权限与架构图渲染。配合 at-least-once、exactly-once、有序投递三种语义以及完备的重试/死信队列机制,你可以在保持代码极简的同时构建高可靠、易扩展的异步系统。作为参考,uptime 教程是使用 Pub/Sub 的事件驱动示例应用,可帮助你快速上手。

【免费下载链接】encoreThe infrastructure platform for the intelligence era项目地址: https://gitcode.com/GitHub_Trending/encor/encore

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/15 12:42:42

不写后端,三步跑起一个类抖音短视频应用

不写后端,三步跑起一个类抖音短视频应用 【免费下载链接】douyin Vue3 Pinia 仿抖音,Vue 在移动端的最佳实践 . Imitate TikTok ,Vue Best practices on Mobile 项目地址: https://gitcode.com/GitHub_Trending/do/douyin Douyin-Vu…

作者头像 李华
网站建设 2026/9/15 12:42:19

用 Twilio/SendGrid 为 IoT 地理围栏触发函数添加短信与邮件通知

用 Twilio/SendGrid 为 IoT 地理围栏触发函数添加短信与邮件通知 【免费下载链接】IoT-For-Beginners 12 Weeks, 24 Lessons, IoT for All! 项目地址: https://gitcode.com/GitHub_Trending/io/IoT-For-Beginners 本文围绕 IoT-For-Beginners 运输项目第 4 课&#xff08…

作者头像 李华