第 5 章 · 数据一致性与消息驱动
本章讲解分布式下的 一致性模型、本地/分布式事务、Saga、事务消息、事件驱动架构(EDA)与幂等设计。
前置:PaaS 消息组件 ch4/ch9;架构第 4 章微服务边界。
5.1 一致性光谱
| 模型 | 说明 | 典型场景 | 延迟窗口 |
|---|
| 强一致 | 读到的永远是最新已提交写 | 账户余额、库存扣减 | 同步完成即一致 |
| 最终一致 | 短暂不一致,稍后对齐 | 搜索索引、报表、推荐 | 秒~分钟 |
| 因果一致 | 有因果关系的事件有序可见 | 订单状态流、评论回复 | 同 key 有序 |
| 读己之所写 | 用户看到自己刚写的数据 | 发帖后立即刷新 | 会话粘滞或主库读 |
CAP 回顾:分区(P)发生时,只能在 一致性 C 与 可用性 A 间权衡;多数互联网业务选 AP + 最终一致,金融核心选 CP。
| 业务 | 推荐模型 | 理由 |
|---|
| 支付扣款 | 强一致 | 不能超扣 |
| 商品搜索 | 最终一致 | 几秒滞后可接受 |
| 订单状态展示 | 因果一致 | 不能先显示「已发货」再「待支付」 |
| 点击量统计 | 最终一致 | 近似即可 |
5.2 本地事务 vs 分布式事务
| 方案 | 适用 | 复杂度 | TPS 影响 | 贤紫组件 |
|---|
| 单库事务 | 单体或未拆库 | 低 | 无 | MySQL InnoDB |
| 2PC / XA | 强一致金融(少用) | 很高 | 显著下降 | — |
| Saga | 长事务、可补偿 | 中 | 中 | 状态机 + MQ |
| TCC | 预留/确认/取消 | 高 | 中 | 自研 + Redis 预留 |
| 事务消息 | 本地写 + 发消息原子 | 中 | 低 | RocketMQ |
| Outbox | 可靠发事件 | 中 | 低 | 定时投递 |
选型决策树:
同一服务同一库? ──是──► 本地事务
│
否
├── 需强一致且步骤少? ──► TCC(慎用)
├── 长流程可补偿? ──► Saga
└── 本地写后要发消息? ──► 事务消息 / Outbox
5.3 Saga 模式(推荐掌握)
5.3.1 编排 vs 协同
| 类型 | 协调者 | 优点 | 缺点 |
|---|
| 编排(Orchestration) | 中央状态机 / order-saga | 流程清晰 | 协调者单点需 HA |
| 协同(Choreography) | 各服务听事件 | 解耦 | 难追踪全局状态 |
创建订单 ──► 预占库存 ──► 创建支付单 ──► 支付成功 ──► 确认库存 ──► 发通知
│ │ │
└── 失败则逆序补偿 ◄─────────┘
cancel release void-pay
| 步骤 | 正向服务 | 正向动作 | 补偿动作 | 幂等键 |
|---|
| 1 | order | create(PENDING) | cancel | orderId |
| 2 | inventory | reserve | release | orderId |
| 3 | payment | create + pay | refund | paymentId |
| 4 | inventory | confirm | — | orderId |
# Saga 状态机片段(概念)
TRANSITIONS = {
"CREATED": {"reserve_ok": "STOCK_RESERVED", "fail": "CANCELLED"},
"STOCK_RESERVED": {"pay_ok": "PAID", "fail": "COMPENSATING"},
"COMPENSATING": {"done": "CANCELLED"},
}
行业案例 · 某跨境电商:用同步 REST 链「下单→扣库存→调支付」,支付超时 30s 导致库存 长期预占。改 Saga + 15min 超时自动 release,预占泄漏下降 95%。
5.4 事务消息(RocketMQ)
┌─────────────┐ 1. 发半消息 ┌──────────┐
│ order-api │ ────────────────► │ RocketMQ │
│ │ 2. 本地事务写库 │ │
│ │ ◄── 3. commit/rollback
└─────────────┘ 4. 消费者可见 └──────────┘
// RocketMQ 事务消息监听器(概念)
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
orderService.createPending(order);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
return orderService.exists(msg.getKeys())
? LocalTransactionState.COMMIT_MESSAGE
: LocalTransactionState.ROLLBACK_MESSAGE;
}
| 要点 | 说明 |
|---|
| 回查 | 网络抖动时 Broker 回查本地事务状态 |
| 幂等 | 消费者按 orderId 去重 |
| 与 Kafka | Kafka 事务更适合流处理;业务事务消息常用 RocketMQ(PaaS ch9) |
5.5 Outbox 模式
-- 同一本地事务内
BEGIN;
INSERT INTO orders (...) VALUES (...);
INSERT INTO outbox (event_type, payload, created_at)
VALUES ('OrderCreated', '{"orderId":"O1"}', NOW());
COMMIT;
-- 独立投递进程轮询 outbox,发 MQ 后标记 sent
| 对比 | 事务消息 | Outbox |
|---|
| MQ 绑定 | 强依赖 Broker 事务 | 任意 MQ |
| 实现 | Broker 原生 | 多一张表 + 投递器 |
| 延迟 | 低 | 轮询间隔(通常 < 1s) |
5.6 事件驱动架构(EDA)
OrderCreated 事件 ──► inventory 订阅(预占)
──► points 订阅(加积分)
──► notification 订阅(发短信)
──► search 订阅(更新索引)
| 优点 | 缺点 | 治理要求 |
|---|
| 解耦 | 调试链路长 | 统一 traceId |
| 易加消费者 | 顺序/重复难 | 幂等 + 分区 key |
| 削峰 | 最终一致窗口 | 监控 lag |
5.6.1 事件规范
{
"eventId": "evt_8f3a2b1c",
"eventType": "OrderCreated",
"occurredAt": "2026-08-18T10:00:00Z",
"aggregateId": "O20260818001",
"version": 1,
"payload": {
"userId": "U100",
"items": [{"sku": "SKU1", "qty": 2}]
},
"traceId": "trace-abc"
}
| 规则 | 说明 |
|---|
| 命名 | 过去式 OrderCreated |
| eventId | 全局唯一,幂等键 |
| version | 载荷演进 |
| traceId | 链路关联 |