返回分布式与数据质量模块支付成功后的异步业务
幂等消费者的原子处理
消费失败的有限处理路径
从业务动作到一致性证据
Distributed Systems / Tutorial 17
消息队列与异步任务测试教程
接受消息传输的不确定性,再用契约、幂等、重试、监控和对账保证业务结果可信。
9 个章节订单事件链重复 丢失 乱序 积压
01
先理解异步链路的三个时钟
返回成功不等于完成支付回调写订单状态
Outbox保存待发事件
Broker持久化与投递
消费者扣库存发积分
查询/对账确认最终状态
交付语义与业务选择
| 语义 | 含义 | 适用判断 |
|---|---|---|
| 至多一次 | 可能丢,不重复 | 非关键通知、可接受遗漏的指标 |
| 至少一次 | 不轻易丢,可能重复 | 订单、库存、支付等常见业务 |
| 效果上的恰好一次 | 传输仍可能重复,业务结果幂等 | 资金、库存和优惠核销 |
HTTP 成功示例统一为 200,但它只说明同步入口已接受请求;支付、库存或退款是否真正完成,要继续检查响应体业务结果和异步最终状态。
02
把消息当作长期演进的业务契约
事件可识别可追踪OrderPaid.v1.json
{
"eventId": "E-1001",
"eventType": "OrderPaid",
"eventVersion": 1,
"occurredAt": "2026-08-10T02:00:00Z",
"aggregateId": "O-1001",
"aggregateVersion": 3,
"traceId": "4bf92f3577b34da6",
"data": {
"paymentId": "P-1001",
"paidAmount": "80.00",
"currency": "CNY"
}
}契约检查
- eventId 全局唯一,用于幂等和审计。
- aggregateId 是分区键,aggregateVersion 判断新旧。
- 时间使用明确时区,金额使用字符串或定点数。
- 新增字段保持向后兼容;删除、改名和改类型需升级版本。
- 消息不携带令牌、卡号或不必要的个人信息。
03
主动制造重复,验证消费者幂等
至少一次投递收到事件读取 eventId
检查去重表是否已处理
业务写入库存或积分
记录已处理同一事务
确认消息提交位点
重复事件测试
def test_当支付事件重复投递时_库存只扣减一次(harness):
event = order_paid(event_id="E-DUP-1", order_id="O-1", quantity=2)
harness.publish("order-events", event)
harness.publish("order-events", event)
harness.wait_consumed("E-DUP-1", attempts=2)
assert harness.inventory("SKU-1")["deducted"] == 2
assert harness.processed_event_count("E-DUP-1") == 1
assert harness.points_ledger_count("O-1") == 1“消费者执行了两次但第二次没报错”不是幂等证据。必须核对库存、优惠核销、积分账本、通知等所有副作用只发生一次。
04
在每个确认窗口测试消息丢失
生产与消费两端丢失窗口与证据
| 故障窗口 | 风险 | 必须验证 |
|---|---|---|
| 生产后进程崩溃 | 事务提交但消息未发 | Outbox 中仍有待投递记录 |
| Broker 确认丢失 | 生产者不知道是否成功 | 同 eventId 重发,消费者幂等 |
| 消费完成前崩溃 | 消息再次投递 | 副作用只发生一次 |
| 消费异常后仍确认 | 消息永久丢失 | 不得 ack,进入重试或死信 |
Outbox 原子写入示意
BEGIN;
UPDATE orders SET status = 'PAID', version = version + 1
WHERE order_id = 'O-1001' AND status = 'PENDING_PAYMENT';
INSERT INTO outbox_events(event_id, aggregate_id, event_type, payload, status)
VALUES ('E-1001', 'O-1001', 'OrderPaid', '{...}', 'PENDING');
COMMIT;
-- 测试:提交后杀死发布进程;重启 relay 后 E-1001 必须最终投递
SELECT status, attempts FROM outbox_events WHERE event_id = 'E-1001';当……时,……关键用例
- 当数据库回滚时,不得留下可投递 Outbox 事件。
- 当 Relay 发送后未更新状态就崩溃时,允许重发但业务结果不重复。
- 当 Broker 暂时不可用时,PENDING 事件保留并带退避重试。
- 当消费者处理失败时,不得提前提交位点。
05
打乱事件顺序,阻止状态倒退
局部有序订单事件乱序测试
| 输入顺序 | 处理策略 |
|---|---|
| PAID(v3) 后收到 CREATED(v1) | 忽略旧版本,订单仍为已支付 |
| REFUNDED(v5) 先于 REFUNDING(v4) | 缓存/延迟或按版本拒绝回退 |
| 同订单事件并行消费 | 按 orderId 分区,保持分区内顺序 |
| 不同订单同时到达 | 允许并行,不建立全局顺序依赖 |
聚合版本保护
def apply_event(order, event):
if event["aggregateVersion"] <= order.version:
return "IGNORED_STALE_EVENT"
if event["aggregateVersion"] != order.version + 1:
return "WAITING_FOR_GAP"
order.transition(event["eventType"])
order.version = event["aggregateVersion"]
return "APPLIED"
def test_当旧事件晚到时_订单状态不回退(order):
order.status, order.version = "PAID", 3
result = apply_event(order, {"eventType": "OrderCreated", "aggregateVersion": 1})
assert result == "IGNORED_STALE_EVENT"
assert order.status == "PAID"06
区分可恢复失败、毒消息与死信
重试有上限首次失败记录原因
指数退避释放下游压力
再次消费保持幂等
超过上限进入 DLQ
修复重放审计结果
重试与死信配置示意
retry:
attempts: 5
backoff: exponential
initialDelayMs: 1000
maxDelayMs: 30000
deadLetter:
topic: order-events.dlq
include:
- originalTopic
- eventId
- failureReason
- attempts
- firstFailedAt分类处理
- 网络超时、依赖限流通常可退避重试。
- Schema 非法、缺必填字段通常直接进入死信。
- 余额不足、退款超额是业务拒绝,不应无限重试。
- DLQ 重放前先修复原因,并用原 eventId 保持幂等。
07
制造积压,验证容量、扩容和恢复
不要只看队列长度积压监控的四个核心指标
| 指标 | 含义 | 判断 |
|---|---|---|
| Consumer lag | 最新位点与消费位点差 | 持续增长说明消费跟不上 |
| Oldest message age | 最老待处理消息年龄 | 直接反映业务延迟 |
| Throughput | 每秒生产/消费量 | 判断扩容是否有效 |
| Retry / DLQ rate | 重试和死信比例 | 发现毒消息或依赖故障 |
积压演练伪命令
# 1. 将消费者缩容到 1 个实例
kubectl scale deployment order-consumer --replicas=1
# 2. 向测试 Topic 发送 100000 条带 runId 的事件
python scripts/publish_orders.py --count 100000 --run-id backlog-01
# 3. 扩容并观察 lag、oldest age、吞吐与错误率
kubectl scale deployment order-consumer --replicas=6
# 4. 验证所有 runId=backlog-01 的业务状态最终收敛只在隔离环境执行容量演练。扩容后不仅要看积压下降,还要检查数据库连接、下游限流、分区分配和消息顺序是否被破坏。
08
用状态收敛和对账证明最终一致
允许延迟,不允许失控支付成功订单 PAID
事件投递OrderPaid
异步消费库存与积分
状态查询等待收敛
对账修复发现并补偿
最终一致性测试
def test_当支付成功事件短暂失败时_业务最终一致(harness):
harness.fail_consumer("inventory", times=2)
response = harness.pay("O-1001")
assert response.status_code == 200
assert response.json()["businessCode"] == "SUCCESS"
harness.wait_until(
lambda: harness.snapshot("O-1001"),
lambda x: x["order"] == "PAID"
and x["inventory"] == "DEDUCTED"
and x["pointsLedger"] == 1,
timeout=30,
)
assert harness.dlq_count(order_id="O-1001") == 0一致性标准必须量化
- 正常情况下多少秒内收敛。
- 依赖恢复后积压多快清空。
- 超时未收敛由什么告警发现。
- 对账任务如何定位、补偿并留下审计记录。
09
完成一次订单事件可靠性演练
练习与验收练习:覆盖异步链路的六类风险
- 为 OrderPaid 定义版本化消息契约和分区键。
- 重复投递同一 eventId,核对库存、积分和通知只执行一次。
- 在事务提交、发送确认和消费确认窗口分别注入崩溃。
- 逆序投递创建、支付和退款事件,验证状态不倒退。
- 模拟可恢复异常与毒消息,核对退避、上限和 DLQ 元数据。
- 修复消费者后按原 eventId 重放死信。
- 制造 10 万条积压并验证扩容、下游容量和恢复时间。
- 定义收敛时限,运行订单、库存、支付、退款对账。
消息可治理
- 契约版本明确
- eventId 唯一
- 分区键正确
- 敏感数据最少
失败可恢复
- 重复不产生副作用
- 丢失窗口覆盖
- 重试有上限
- 死信可审计重放
结果可证明
- 乱序不倒退
- 积压有指标
- 收敛时限明确
- 对账补偿可执行
你已经能验证异步消息的可靠性。下一步聚焦缓存故障、热点与数据库一致性。
继续学习缓存测试