第7章:Pub/Sub —— 主题、订阅与事件驱动架构
约 4 分钟 · 更新于 2026-09-02
同步调用把两个服务的可用性绑在一起;事件把它们解开。Encore 的 Pub/Sub 原语用两个类(Topic / Subscription)覆盖发布订阅全场景,本地 NSQ、云上 SNS+SQS / GCP Pub/Sub,代码零改动。
一、为什么需要 Pub/Sub
以订票场景为例:订单支付成功后要发确认短信、更新库存、通知报表系统。同步实现的问题:
- 支付接口的响应时间 = 自身 + 三个下游之和;
- 任何一个下游挂了,支付"失败";
- 每加一个下游都要改支付代码。
事件驱动的版本:支付服务只发布 order-paid 事件;短信、库存、报表各自订阅。发布者不知道(也不关心)有谁在听——新增消费者不改发布方代码。
二、声明主题 (Topic)
typescript
import { Topic } from "encore.dev/pubsub";
export interface OrderPaidEvent {
orderNo: string;
amountCents: number;
paidAt: string;
}
export const orderPaid = new Topic<OrderPaidEvent>("order-paid", {
deliveryGuarantee: "at-least-once",
});
规则:
- 主题必须是包级变量(静态分析要求),不能在函数里创建;
- 事件类型是普通接口,作为消息 schema 同时用于类型检查与文档;
- 主题从任何服务都可访问(export 后 import 即可)——但发布权通常约定归属一个领域服务。
三、发布与订阅
typescript
// 发布(任何服务内)
const messageID = await orderPaid.publish({
orderNo: "ORD123",
amountCents: 158800,
paidAt: new Date().toISOString(),
});
typescript
// 订阅(消费方服务内)
import { Subscription } from "encore.dev/pubsub";
import { orderPaid } from "../order/events";
const _ = new Subscription(orderPaid, "send-confirm-sms", {
handler: async (event) => {
await sendSms(event.orderNo);
},
});
- 订阅名("send-confirm-sms")在主题内唯一,标识一个独立消费组;
- 同一主题可挂任意多个订阅,各自独立收到全量消息;
- handler 抛异常 = 消费失败,触发重试。
四、投递语义
| 语义 | 行为 | 代价与限制 |
|---|
| at-least-once(默认) | 至少投递一次,可能重复 | handler 必须幂等 |
| exactly-once | 最小化重复投递 | 云上限流:AWS 300 msg/s/主题,GCP 3000+ msg/s/区域;且不做发布端去重 |
幂等是订阅方的义务
at-least-once 下重复投递是正常现象。幂等三板斧:
- 数据库唯一约束(如 sms_log.order_no 唯一,重复插入吞掉 already_exists);
- 状态机前置检查(已是 NOTIFIED 就直接 return);
- 用 currentRequest() 的 deliveryAttempt 感知重试次数,超阈值走补偿逻辑。
即使 exactly-once 也不能豁免——它只是"最小化"重复,且发布端重试仍可能产生重复消息。
失败与死信
投递失败按重试策略退避重试;超过最大次数进入死信队列 (Dead Letter Queue, DLQ),不阻塞后续消息。DLQ 消息在云上对应 SQS DLQ / GCP dead-letter topic,需要运维流程兜底(告警 + 人工/脚本重放)。
五、消息属性与有序投递
字段包上 Attribute<T> 后成为消息属性(不在 body 里,可用于过滤/排序键):
typescript
import { Topic, Attribute } from "encore.dev/pubsub";
export interface CartEvent {
shoppingCartID: Attribute<number>; // 属性字段
event: string;
}
export const cartEvents = new Topic<CartEvent>("cart-events", {
deliveryGuarantee: "at-least-once",
orderingAttribute: "shoppingCartID", // 同一购物车的事件按序投递
});
orderingAttribute 保证同一键值内的消息按发布顺序投递(不同键值之间无序)。限制:AWS 300 msg/s/主题、GCP 1 MBps/排序键;本地环境不模拟有序,别在本地"验证"顺序性。
六、本地与云的映射
| 环境 | 实现 |
|---|
| 本地 encore run | NSQ(Docker 自动拉起) |
| AWS | SNS + SQS |
| GCP | Pub/Sub |
代码不变。这就是原语抽象的价值:语义在框架层锁定,实现按环境替换。
七、事件驱动的架构模式
模式一:状态翻转才发事件(第 14 章实战)
轮询类系统(监控、对账)里"每次检查都发事件"会淹没消费者。正确做法:对比上一次状态,翻转时才发布:
typescript
const wasUp = await getPreviousMeasurement(site.id);
if (up !== wasUp) {
await TransitionTopic.publish({ site, up });
}
模式二:事务性发件箱的替代
"写库 + 发事件"的原子性问题(写库成功、发布失败)在 Encore 里同样存在。轻量解法:先写库、后发布,订阅方幂等 + 发布方失败重试;严格场景用 outbox 表 + Cron 扫描补发(第 8 章的 CronJob 正好派上用场)。
模式三:扇出编排
一个业务动作触发多个后续步骤时,用一个事件 + N 个订阅,而不是在业务代码里排队调 N 个服务。第 14 章的 uptime 监控就是最小案例:uptime-transition 主题 + Slack 通知订阅。
八、测试 Pub/Sub
encore test 环境下发布是真实可调用的(本地实现),常规做法:
typescript
import { describe, it, expect } from "vitest";
import { orderPaid } from "./events";
describe("pub/sub", () => {
it("publishes order paid event", async () => {
const messageId = await orderPaid.publish({
orderNo: "ORD123",
amountCents: 158800,
paidAt: new Date().toISOString(),
});
expect(messageId).toBeDefined();
});
});
订阅 handler 是普通函数——直接调用它测试业务逻辑,不必绕道真实投递链路。投递语义(重试、DLQ)是框架职责,不需要你的测试覆盖。
九、常见误区
- handler 不幂等——at-least-once 下必然重复消费,迟早出重复短信/重复入账。
- 把 Pub/Sub 当任务队列精确控制消费节奏——它是广播语义;需要工作队列语义时用"单订阅 + 幂等"逼近,或另选专门组件。
- 依赖跨键全局有序——orderingAttribute 只保证同键内有序。
- 在本地验证 exactly-once / 有序性——本地 NSQ 不模拟这些语义。
- 用事件传"胖数据"(整个订单对象)——事件传引用(orderNo)+ 消费方回查,除非明确要做事件溯源。
- 发布失败没有兜底——严格一致场景上 outbox 模式。
- 在函数体内 new Topic——包级变量硬规则。
十、实践练习
- 给第 6 章的订单应用加 order-paid 主题:支付端点发布事件,notification 服务订阅并打日志。
- 在 handler 里故意抛异常,观察本地重试行为(看仪表盘 trace 里的重复投递)。
- 给通知落一张 notify_log 表加唯一约束,实现幂等消费,再次制造重试验证不重复发送。
- 加第二个订阅 update-report,体会"新增消费者不改发布方"。
- 用 Attribute + orderingAttribute 声明一个按订单号有序的主题,写出它的适用场景说明(提示:同一订单的状态流转事件)。
十一、总结
- Pub/Sub 解耦服务可用性与响应时间:发布方不知道消费方存在。
- Topic 包级声明 + 类型化事件接口;Subscription 独立消费组,多订阅各收全量。
- 投递语义默认 at-least-once,幂等是订阅方义务;exactly-once 有限流且不豁免幂等。
- 失败重试 → DLQ,需运维兜底;orderingAttribute 提供同键有序。
- 本地 NSQ / AWS SNS+SQS / GCP Pub/Sub 同一套代码。
- 三个实战模式:状态翻转才发布、outbox 补发、事件扇出。
请继续阅读:第8章:定时任务与密钥。
原始资料引用