资讯动态

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

发布时间:2026/9/15 12:44:29 来源:尧图企业网站定制
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 调用调用链如下user服务开启数据库事务写入用户记录user服务调用email服务发送欢迎邮件email服务再调用邮件供应商真正发送邮件邮件发送成功后email服务回复user服务user服务再调用analytics服务记录注册信息analytics服务写入数据仓库analytics服务回复user服务user服务提交数据库事务user服务才回复用户“注册成功”。这种设计有两个明显问题响应时间被下游拖累——如果邮件供应商耗时 3 秒用户就要等 3 秒才能收到注册成功响应而实际上用户写入数据库那一刻就可以确认注册故障爆炸半径大——如果数据仓库故障所有用户注册都会失败而分析系统是纯内部功能不应该影响用户注册。Pub/Sub 方式的优势改用 Pub/Sub 后注册流程变为user服务开启数据库事务写入用户记录向signupsTopic 发布一条注册事件提交事务立即回复用户注册成功。此时user服务与email、analytics服务彻底解耦后两者各自订阅signupsTopic 并行处理事件若处理失败事件会自动退避重试直到成功或达到最大重试次数后进入死信队列DLQ。user服务甚至完全不知道email、analytics的存在未来新增任何关注“新用户注册”的系统都无需改动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*SignupEventTopicConfig 配置项从 runtimes/go/pubsub/internal/types/public.go 的定义看TopicConfig包含两个字段字段类型必填说明DeliveryGuaranteeDeliveryGuarantee是投递语义AtLeastOnce或ExactlyOnceOrderingAttributestring否作为排序键的消息属性名设置后同一键值的消息按发布顺序投递运行时topic.go在创建 Topic 时会根据当前运行环境选择实现单元测试场景使用internal/test的测试实现未注册的场景降级为 noop 实现正常部署时则按PubsubProviders匹配对应的云厂商实现仓库内置 AWS、GCP、Azure、Encore Cloud 与本地 NSQ 五种 provider。投递语义从 At-Least-Once 到 Exactly-OnceAt-least-once至少一次投递上面的示例配置保证对于每个订阅事件至少被投递一次。如果 Topic 认为事件未被成功处理会尝试再次投递。因此所有订阅处理函数都应设计为幂等的——即处理函数被调用两次或多次时从外部看与调用一次没有差别。实现幂等通常有两种方式用数据库记录该事件触发的动作是否已执行或者确保动作本身天然具备幂等性。Exactly-once精确一次投递将DeliveryGuarantee设置为pubsub.ExactlyOnce可在基础设施层面提供更强的保证最小化消息被重复投递的可能性var Signups pubsub.NewTopic*SignupEvent但即便如此仍有极少数情况下消息会被重投例如网络问题导致“处理成功”的确认消息在到达云厂商之前丢失即著名的两军问题 也明确指出Exactly-once 只约束“投递给消费者”这一环节不包含发布侧的去重——如果应用逻辑中Publish被调用了两次例如应用层重试消息会被投递两次且 Exactly-once Topic 上的订阅相比 At-least-once 有更高的投递延迟。启用 exactly-once 后云厂商会施加吞吐限制AWSTopic 每秒最多 300 条消息参见 AWS SQS QuotasGCP区域范围内所有 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.goPublish会先从消息中提取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调用背后包含完整的处理链路上下文与合法性检查context 已取消则直接返回Topic 未通过NewTopic创建则返回Unimplemented错误属性提取与 JSON 序列化通过utils.MarshalFields提取pubsub-attr标记的属性并将消息体序列化为 JSON序列化失败返回InvalidArgument排序键解析若配置了OrderingAttribute从属性中取出排序键追踪与关联 ID 传播若当前处于某个请求中会向属性注入encore_parent_trace_id、encore_ext_correlation_id外部关联 ID 优先否则用 trace ID以及平台请求的强制追踪标记让订阅者可以把 trace 标记为当前请求的子 span——这正是 Encore 分布式追踪在异步消息链路上保持连贯的实现机制限流与发布等待发布限流器publishLimiter放行后调用对应云厂商的PublishMessage真正发布错误归一化发布失败统一包装为Unavailable错误码返回。注意Publish会阻塞直到消息被 Topic 成功接收若返回错误通常意味着发布失败但不排除消息仍可能被订阅者收到分布式系统的经典语义。使用 TopicRef突破静态分析的引用限制Encore 通过静态分析确定哪些服务向哪些 Topic 发布消息据此正确配置基础设施、渲染架构图并配置 IAM 权限。因此*pubsub.Topic变量不能随意传递——那会使静态分析在许多场景下失效。为此 Encore 提供了TopicRef它返回一个可以自由传递的 Topic 引用signupRef : pubsub.TopicRef[pubsub.Publisher[*SignupEvent]](Signups) // signupRef 的类型是 pubsub.Publisher[*SignupEvent]仅允许发布操作TopicRef与Topic的关键区别是引用必须预先声明所需权限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 时使用默认值 1000 时表示重试 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 发布不在同一事务中时可能导致服务间数据不一致。在不引入过多复杂度的前提下推荐采用事务性 Outboxtransactional 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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价