返回分布式与数据质量模块
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

完成一次订单事件可靠性演练

练习与验收

练习:覆盖异步链路的六类风险

  1. 为 OrderPaid 定义版本化消息契约和分区键。
  2. 重复投递同一 eventId,核对库存、积分和通知只执行一次。
  3. 在事务提交、发送确认和消费确认窗口分别注入崩溃。
  4. 逆序投递创建、支付和退款事件,验证状态不倒退。
  5. 模拟可恢复异常与毒消息,核对退避、上限和 DLQ 元数据。
  6. 修复消费者后按原 eventId 重放死信。
  7. 制造 10 万条积压并验证扩容、下游容量和恢复时间。
  8. 定义收敛时限,运行订单、库存、支付、退款对账。

消息可治理

  • 契约版本明确
  • eventId 唯一
  • 分区键正确
  • 敏感数据最少

失败可恢复

  • 重复不产生副作用
  • 丢失窗口覆盖
  • 重试有上限
  • 死信可审计重放

结果可证明

  • 乱序不倒退
  • 积压有指标
  • 收敛时限明确
  • 对账补偿可执行

你已经能验证异步消息的可靠性。下一步聚焦缓存故障、热点与数据库一致性。

继续学习缓存测试