GraphQL 的查询与变更解决的是"拉"的问题,但真实系统里大量场景需要"推":订单状态变了要通知客户端、支付成功了要回调商户、库存变动要触发下游。这些"推"的需求对应三种范式——Subscription(面向已连接的客户端)、Webhook(面向外部系统)、消息队列(面向内部服务)。三者看起来都是"发事件",但投递语义、可靠性、适用场景截然不同。本文系统梳理这三种范式的定位、事件 Schema 设计、投递保证与重试、幂等与顺序、死信与可观测性。订阅的实现细节可先阅读 https://plumephp.com/graphql-realtime-subscription-sse/;错误处理与重试可参考 https://plumephp.com/graphql-error-handling/。
一、三种异步范式的定位
1.1 一张图看清差异
| 维度 | Subscription | Webhook | 消息队列 |
|---|---|---|---|
| 消费者 | 已连接的客户端 | 外部系统 | 内部服务 |
| 连接方向 | 长连接(双向) | 服务端主动 POST | 生产者/消费者解耦 |
| 投递保证 | 尽力而为 | 至少一次 + 重试 | 取决于 MQ |
| 顺序保证 | 连接内有序 | 无 | 分区内有序 |
| 典型场景 | 实时 UI 更新 | 商户回调、集成 | 领域事件、异步任务 |
| 失败处理 | 客户端重连 | 重试 + DLQ | 重试 + DLQ |
1.2 选型的三个问题
- 谁在消费:自家客户端 → Subscription;外部系统 → Webhook;自家服务 → MQ。
- 能否容忍丢消息:UI 更新可容忍 → Subscription;资金相关不可容忍 → MQ。
- 是否需要解耦:需要削峰/解耦 → MQ;只是通知 → Webhook。
1.3 三者可以组合
领域事件(MQ)──┬──► 内部服务消费(库存、风控)
├──► Subscription 推送(在线用户实时 UI)
└──► Webhook 投递(外部商户回调)
同一个领域事件,可以同时驱动三种出口——这是事件驱动架构的核心价值。
一句话总结:Subscription 面向"在线的自己人",Webhook 面向"外部系统",消息队列面向"内部服务"——先确定消费者是谁,范式就定了。
二、事件 Schema 设计
2.1 事件也是 Schema,也要治理
事件一旦发布就是契约,消费方依赖它。事件 Schema 的变更管理与 GraphQL Schema 同等重要。
# 事件类型定义(可作为 GraphQL 类型暴露,也可映射到 Avro/Protobuf)
type OrderCreatedEvent {
eventId: ID! # 全局唯一,用于幂等
occurredAt: DateTime! # 事件发生时间(不是投递时间)
aggregateId: ID! # 聚合根 ID(订单 ID)
version: Int! # 事件版本,用于演进
payload: OrderCreatedPayload!
}
type OrderCreatedPayload {
orderId: ID!
buyerId: ID!
total: Money!
items: [OrderItemSnapshot!]!
}
2.2 事件设计的五条原则
| 原则 | 说明 |
|---|---|
| 不可变 | 事件是事实,发布后不可修改 |
| 自包含 | 携带消费所需的最小充分信息 |
| 有版本 | version 字段支持演进 |
| 有标识 | eventId 支持幂等去重 |
| 有语义 | 用过去时命名(OrderCreated 而非 CreateOrder) |
2.3 事件信封与载荷分离
// 统一信封:元数据与业务载荷分离
interface EventEnvelope<T> {
eventId: string;
eventType: string; // "order.created"
eventVersion: number;
occurredAt: string; // ISO 8601
traceId: string; // 全链路追踪
producer: string; // 生产服务名
payload: T;
}
2.4 事件命名规范
# 事件类型命名:<domain>.<aggregate>.<action>
order.created # 订单已创建
order.paid # 订单已支付
payment.refunded # 支付已退款
inventory.reserved # 库存已预留
# ❌ 避免命令式命名
create_order # 这是命令,不是事件
一句话总结:事件是"已经发生的事实",用过去时命名、带唯一 ID、带版本、自包含——这四条决定了事件能否被安全消费与演进。
三、订阅:面向客户端的实时推送
3.1 订阅的执行模型
Subscription 建立在长连接(WebSocket / SSE)之上,每个订阅对应一个"事件流 + 过滤器"。
# 客户端订阅:只关心自己订单的状态变化
subscription OnOrderStatus($orderId: ID!) {
orderStatusChanged(orderId: $orderId) {
orderId
status
updatedAt
}
}
3.2 服务端实现
import { PubSub } from 'graphql-subscriptions';
const pubsub = new PubSub();
const resolvers = {
Subscription: {
orderStatusChanged: {
subscribe: (_p, { orderId }, ctx) => {
// 鉴权:确认用户有权订阅该订单
assertCanViewOrder(ctx.user, orderId);
// 过滤:只推送该订单的事件
return pubsub.asyncIterator(`order.status.${orderId}`);
},
},
},
Mutation: {
updateOrderStatus: async (_p, { id, status }, ctx) => {
const order = await orderService.update(ctx, id, status);
// 发布事件到对应频道
await pubsub.publish(`order.status.${id}`, {
orderStatusChanged: { orderId: id, status, updatedAt: new Date() },
});
return order;
},
},
};
3.3 生产环境的关键约束
| 约束 | 说明 |
|---|---|
| 无状态限制 | PubSub 内存实现不支持多实例,需 Redis / NATS |
| 鉴权 | 订阅必须在 subscribe 中鉴权,不能只靠过滤器 |
| 背压 | 客户端消费慢时要丢弃或断开 |
| 断线重连 | 客户端需实现指数退避重连 |
| 扩容 | 长连接数决定实例数与连接层设计 |
// 多实例:用 Redis PubSub 替换内存实现
import { RedisPubSub } from 'graphql-redis-subscriptions';
const pubsub = new RedisPubSub({
publisher: new Redis(process.env.REDIS_URL),
subscriber: new Redis(process.env.REDIS_URL),
});
一句话总结:Subscription 是"尽力而为"的实时通道——它能提升体验,但不能承担"必须送达"的业务语义,资金类通知必须走 MQ 或 Webhook。
四、Webhook:面向外部系统的投递
4.1 Webhook 的本质
Webhook 是"服务端主动发起的 HTTP POST",把事件推送到订阅方提供的 URL。它的可靠性完全依赖重试 + 幂等 + 签名三件套。
4.2 投递流程
事件产生 → 入队 → 投递器 → HTTP POST → 目标系统
2xx → 标记成功;非 2xx → 退避重试 → 超过上限 → 死信队列
4.3 签名与验签
import { createHmac } from 'node:crypto';
// 发送方:对 (timestamp + body) 签名
function sign(body: string, timestamp: number, secret: string): string {
return createHmac('sha256', secret)
.update(`${timestamp}.${body}`)
.digest('hex');
}
// 投递时附带签名头
await fetch(endpoint, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'X-Webhook-Id': event.eventId,
'X-Webhook-Timestamp': String(Date.now()),
'X-Webhook-Signature': sign(JSON.stringify(event), Date.now(), secret),
},
body: JSON.stringify(event),
});
// 接收方:验签 + 防重放(时间窗口)
function verify(raw: string, sig: string, ts: string, secret: string): boolean {
if (Date.now() - Number(ts) > 5 * 60 * 1000) return false; // 超过 5 分钟拒绝
return timingSafeEqual(Buffer.from(sig), Buffer.from(sign(raw, Number(ts), secret)));
}
4.4 Webhook 的可靠性设计
| 设计点 | 做法 |
|---|---|
| 至少一次投递 | 2xx 才算成功,否则重试 |
| 指数退避 | 1s → 5s → 30s → 5m → 30m |
| 最大重试 | 8~12 次后入 DLQ |
| 超时 | 连接 5s、读取 10s |
| 并发控制 | 同一 endpoint 串行或限并发 |
| 幂等 | 携带 eventId,接收方去重 |
一句话总结:Webhook 的可靠性不靠"保证不丢",而靠"丢了能重试、重了能去重、被伪造能验签"——三件套缺一不可。
五、消息队列:内部服务的解耦
5.1 为什么内部用 MQ 而非 Webhook
| 需求 | MQ | Webhook |
|---|---|---|
| 削峰 | 天然支持 | 不支持 |
| 多消费者 | 广播/组播 | 一对多需自行管理 |
| 顺序 | 分区内有序 | 无 |
| 回溯重放 | 支持 | 不支持 |
| 解耦 | 强 | 弱(依赖对方可用) |
5.2 事件发布
// 事务性发件箱(Transactional Outbox):避免"写库成功但发消息失败"
async function createOrder(tx: Tx, input: CreateOrderInput) {
const order = await tx.order.create({ data: input });
// 事件与业务写入同一事务,保证原子性
await tx.outbox.create({
data: {
eventId: randomUUID(),
eventType: 'order.created',
aggregateId: order.id,
payload: JSON.stringify(toEventPayload(order)),
status: 'PENDING',
},
});
return order;
}
独立的 relay 进程轮询 outbox 表中的 PENDING 记录,投递到 MQ 后置为 SENT,从而在"业务写入"与"事件发布"之间建立可靠桥梁。
5.3 分区与顺序
用 aggregateId 作为分区键(partition key),可保证同一聚合的事件进入同一分区,从而获得分区内有序性。
一句话总结:内部事件用 MQ,关键是用事务性发件箱解决"业务写入与事件发布"的原子性——直接发消息是最常见的丢事件源头。
六、投递保证与重试
6.1 三种投递语义
| 语义 | 含义 | 代价 |
|---|---|---|
| 至多一次 | 可能丢,不重复 | 低 |
| 至少一次 | 不丢,可能重复 | 需幂等 |
| 恰好一次 | 不丢不重 | 高(端到端很难真正实现) |
工程上最务实的选择是"至少一次 + 幂等",而不是追求端到端的"恰好一次"。
6.2 重试策略
// 指数退避 + 抖动,避免重试风暴
function backoff(attempt: number): number {
const base = Math.min(2 ** attempt * 1000, 5 * 60 * 1000); // 上限 5 分钟
const jitter = Math.random() * 1000;
return base + jitter;
}
async function deliverWithRetry(job: DeliveryJob) {
for (let attempt = 1; attempt <= MAX_ATTEMPTS; attempt++) {
try {
const res = await postWithTimeout(job.endpoint, job.body, 10_000);
if (res.ok) return { ok: true };
if (res.status >= 400 && res.status < 500 && res.status !== 429) {
return { ok: false, permanent: true }; // 4xx(除 429)不重试
}
} catch (err) {
logger.warn({ err, attempt }, 'delivery failed');
}
await sleep(backoff(attempt));
}
return { ok: false, permanent: false }; // 转入 DLQ
}
6.3 哪些错误该重试
| 响应 | 是否重试 | 理由 |
|---|---|---|
| 2xx | 不重试(成功) | — |
| 429 | 重试(更长退避) | 限流,稍后可行 |
| 500 / 502 / 503 | 重试 | 临时故障 |
| 400 / 422 | 不重试 | 请求本身有问题 |
| 401 / 403 | 不重试(告警) | 凭证失效,需人工 |
| 超时 / 连接失败 | 重试 | 网络抖动 |
一句话总结:重试必须区分"临时故障"与"永久失败"——对 4xx 盲目重试只会浪费资源并掩盖真正的配置问题。
七、幂等与顺序
7.1 幂等是"至少一次"的必要配套
// 消费端幂等:用 eventId 去重表
async function handleEvent(event: EventEnvelope<unknown>) {
const inserted = await prisma.processedEvent.createMany({
data: [{ eventId: event.eventId }],
skipDuplicates: true, // 唯一约束冲突则跳过
});
if (inserted.count === 0) {
logger.info({ eventId: event.eventId }, 'duplicate event ignored');
return; // 已处理过,直接返回
}
await applyEvent(event);
}
7.2 幂等的三种实现
| 方式 | 适用 | 注意 |
|---|---|---|
| 去重表 | 通用 | 需与业务操作同事务 |
| 状态机 | 有明确状态流转 | 幂等 = 状态不倒退 |
| 版本号 | 有单调版本 | 旧版本丢弃 |
// 状态机幂等:只在合法状态转移时应用
const VALID: Record<string, string[]> = {
CREATED: ['PAID', 'CANCELLED'],
PAID: ['SHIPPED', 'REFUNDED'],
SHIPPED: ['DELIVERED'],
};
function canTransition(from: string, to: string): boolean {
return VALID[from]?.includes(to) ?? false;
}
7.3 顺序问题
需要顺序的场景(同一聚合的事件必须有序),保证手段有三:MQ 分区键取 aggregateId、消费端按分区单线程处理、事件携带 version 以便检测乱序。
// 版本检测:拒绝乱序的旧事件
if (event.version <= current.version) {
logger.warn({ eventId: event.eventId }, 'out-of-order event dropped');
return;
}
一句话总结:幂等解决"重复",版本号解决"乱序"——两者是事件驱动系统稳定运行的底线,缺一个都会在流量上来后爆发。
八、死信队列与可观测性
8.1 死信队列(DLQ)
当事件超过最大重试次数仍未成功,不能直接丢弃,必须进入 DLQ 供人工处理。
// 投递失败 → 写入 DLQ(含失败原因与完整上下文)
async function moveToDlq(job: DeliveryJob, reason: string) {
await prisma.deadLetter.create({
data: {
eventId: job.eventId,
endpoint: job.endpoint,
payload: job.body,
attempts: job.attempts,
lastError: reason,
failedAt: new Date(),
},
});
metrics.dlqSize.inc();
alerts.warn(`Event ${job.eventId} moved to DLQ`);
}
8.2 DLQ 的处理流程
| 步骤 | 动作 |
|---|---|
| 告警 | DLQ 非空即告警 |
| 分类 | 按失败原因分组(网络 / 4xx / 5xx) |
| 修复 | 修复配置或代码 |
| 重放 | 从 DLQ 重新投递 |
| 归档 | 超过保留期后归档 |
8.3 可观测性指标
| 指标 | 含义 | 告警 |
|---|---|---|
| 投递成功率 | 成功 / 总投递 | < 99% |
| 投递延迟 P95 | 事件发生到送达 | > 30s |
| 重试率 | 需重试的比例 | 突增 |
| DLQ 深度 | 死信堆积量 | > 0 |
| 重复率 | 幂等命中比例 | 异常升高 |
| 积压(lag) | MQ 消费滞后 | 持续增长 |
// 用 OpenTelemetry 把事件投递纳入全链路追踪
const span = tracer.startSpan('webhook.deliver', {
attributes: { 'event.id': event.eventId, 'delivery.attempt': job.attempts },
});
事件驱动的可观测性要求"从事件产生到被消费"的全链路可见——traceId 必须贯穿生产、队列、投递与消费四个环节,否则排障时只能靠猜。指标、日志、追踪三支柱在此场景下缺一不可。
Subscription、Webhook、消息队列不是三选一的竞争关系,而是三个互补的出口。事件驱动的成熟度,体现在你是否回答了这些问题:事件有唯一 ID 吗?消费端幂等吗?乱序怎么办?重试区分了临时与永久吗?失败进 DLQ 了吗?全链路能追踪吗?把这六个问题答完,事件驱动才算真正落地。
一句话总结
事件驱动集成的核心是"一事件多出口 + 至少一次投递 + 幂等消费 + 版本防乱序 + DLQ 兜底 + 全链路追踪"——六者齐备,事件才是资产而非负担。
FAQ
Q1: Subscription 能替代 Webhook 吗?
A: 不能。Subscription 依赖客户端保持长连接且在线,适合实时 UI;Webhook 是服务端主动投递,适合外部系统与离线场景。两者的可靠性模型完全不同,不能互相替代。
Q2: 为什么推荐"事务性发件箱"而不是直接发消息?
A: 因为"写数据库"和"发消息"是两个系统,无法用单个事务覆盖。直接发消息会出现"库写成功但消息没发出去"(丢事件)或"消息发了但库回滚"(幽灵事件)。发件箱把事件与业务写入放在同一事务,再由 relay 可靠投递。
Q3: “恰好一次"投递真的做不到吗?
A: 端到端的恰好一次很难实现(需要跨生产、队列、消费三方协调)。务实做法是"至少一次投递 + 消费端幂等”,其效果等价于恰好一次,且实现简单、可验证。
Q4: Webhook 接收方一直返回 500 怎么办?
A: 重试到上限后转入 DLQ 并告警,同时联系接收方排查。不要无限重试(会拖垮投递器),也不要静默丢弃(会丢业务事件)。DLQ 是"延迟处理"而非"放弃"。
Q5: 事件 Schema 变更如何处理兼容性?
A: 遵循与 GraphQL Schema 相同的原则:只增不改,新增字段保持可选,用 eventVersion 标识结构版本,消费端对未知字段宽容(前向兼容)。破坏性变更需新增事件类型而非修改旧类型。
相关阅读
- https://plumephp.com/graphql-observability-tracing/ —— 全链路追踪与指标采集
- Kafka 专题 —— 消息队列的分区、顺序与消费语义
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。