LabHub
学习 学习路径 课程

微服务架构

用 Outbox 模式保证事件发布

在 LabHub 中继续学习

目标

先亲手制造双重写入的不一致,再用发件箱模式消除丢失,并通过消费者幂等性吸收剩余的重复,亲手完成整个链路。

为什么重要

并列调用数据库保存和事件发布的代码无处不在,而且平时运行良好。问题只会在代理恰好波动 3 秒的那一刻出现,当时产生的不一致会悄无声息地残留下来。比如几天后对账时,才发现订单存在但库存没有减少。发件箱模式不采用“让代理参与事务”的方案,而是将思路转变为“把发布意图一并写入数据库”。但这也会带来新的性质——不会再丢失,却会产生重复。如果中继在发布后、更新状态前宕机,重启后就会再次发送同一个事件。因此,本实验将重复视为设计前提而非缺陷,并在同一个实验中一直做到由消费者侧吸收重复。

步骤

  1. 使用 /root/outbox/init.py 创建 /root/outbox/app.db。必须包含 orders(id, sku, qty)outbox(event_id, aggregate_id, seq, event_type, payload, status) 两张表,且 status 的默认值为 PENDING
  2. /root/outbox/dualwrite.py 保存一条订单后发布失败。执行后,在 /root/outbox/dualwrite.out 的第一行写入 INCONSISTENT orders=<n> published=<m>。n 和 m 必须不同。
  3. /root/outbox/place_order.py 在同一个事务中写入 orders 和 outbox。执行 3 次后,orders 应有 3 行,outbox 应有 3 行。
  4. /root/outbox/relay.pystatus='PENDING' 的行 RPUSH 到 Redis 列表 outbox.events,并把该行改为 PUBLISHED。即使执行两次,列表长度也不应增加。
  5. 同一个 aggregate_id 的事件必须按 seq 升序进入队列。在 /root/outbox/order_check.out 中留下 ORDER OK
  6. /root/outbox/relay_crash.py 只发布而不更新状态。执行后再次运行正常中继,同一个 event_id 会进入队列 2 次。在 /root/outbox/atleastonce.out 中写入 DUPLICATE event_id=<id> count=2
  7. /root/outbox/consumer.py 清空队列,同时以 event_id 为基准过滤重复并进行处理。将处理结果保存在 Redis 哈希 processed 中。即使存在重复,processed 的大小也必须等于唯一事件数。
  8. /root/outbox/report.txt 中写入 orders=<n>outbox=<n>enqueued=<n>processed_unique=<n> 四行。enqueued 必须大于或等于 processed_unique

参考

创建订单表和发件箱表

使用 /root/outbox/init.py 创建 /root/outbox/app.db。必须包含 orders(id, sku, qty)outbox(event_id, aggregate_id, seq, event_type, payload, status) 两张表,且 status 的默认值为 PENDING

使用 python3 的 sqlite3 模块就足够了。发件箱行中需要包含事件标识符、聚合标识符、类型、正文和状态。

复现双重写入的不一致

/root/outbox/dualwrite.py 保存一条订单后发布失败。执行后,在 /root/outbox/dualwrite.out 的第一行写入 INCONSISTENT orders=<n> published=<m>。n 和 m 必须不同。

让数据库保存成功,只让代理发布失败即可。统计两个存储中的数量,并把两者不同这一事实记录到文件中。

合并到一个事务中

/root/outbox/place_order.py 在同一个事务中写入 orders 和 outbox。执行 3 次后,orders 应有 3 行,outbox 应有 3 行。

将两个 INSERT 放在同一个连接、同一次提交中。如果提交前发生异常,两者都不存在才是正常结果。

通过中继转移到代理

/root/outbox/relay.pystatus='PENDING' 的行 RPUSH 到 Redis 列表 outbox.events,并把该行改为 PUBLISHED。即使执行两次,列表长度也不应增加。

读取 PENDING 行,将其推入 Redis 列表,然后更改状态。即使多次执行,也不应再次转移已经转移过的内容。

保持每个聚合的顺序

同一个 aggregate_id 的事件必须按 seq 升序进入队列。在 /root/outbox/order_check.out 中留下 ORDER OK

同一个订单的事件必须按照发生顺序发出。思考应采用什么作为排序依据——时间可能相同。

通过中继崩溃制造重复

/root/outbox/relay_crash.py 只发布而不更新状态。执行后再次运行正常中继,同一个 event_id 会进入队列 2 次。在 /root/outbox/atleastonce.out 中写入 DUPLICATE event_id=<id> count=2

模拟已经发布却未能更新状态就宕机的情况。再次运行时,同一个事件会两次进入队列。

为消费者添加去重

/root/outbox/consumer.py 清空队列,同时以 event_id 为基准过滤重复并进行处理。将处理结果保存在 Redis 哈希 processed 中。即使存在重复,processed 的大小也必须等于唯一事件数。

记住已经处理过的事件标识符即可。Redis 的集合数据结构或 SET NX 很合适。

用全链路指标生成报告

/root/outbox/report.txt 中写入 orders=<n>outbox=<n>enqueued=<n>processed_unique=<n> 四行。enqueued 必须大于或等于 processed_unique

将订单数、发件箱行数、进入队列的数量以及实际处理的唯一数量汇总到一个文件中。前三个值与最后一个值之间的关系就是该模式的全部要义。