用 Outbox 模式保证事件发布
目标
先亲手制造双重写入的不一致,再用发件箱模式消除丢失,并通过消费者幂等性吸收剩余的重复,亲手完成整个链路。
为什么重要
并列调用数据库保存和事件发布的代码无处不在,而且平时运行良好。问题只会在代理恰好波动 3 秒的那一刻出现,当时产生的不一致会悄无声息地残留下来。比如几天后对账时,才发现订单存在但库存没有减少。发件箱模式不采用“让代理参与事务”的方案,而是将思路转变为“把发布意图一并写入数据库”。但这也会带来新的性质——不会再丢失,却会产生重复。如果中继在发布后、更新状态前宕机,重启后就会再次发送同一个事件。因此,本实验将重复视为设计前提而非缺陷,并在同一个实验中一直做到由消费者侧吸收重复。
步骤
- 使用
/root/outbox/init.py创建/root/outbox/app.db。必须包含orders(id, sku, qty)和outbox(event_id, aggregate_id, seq, event_type, payload, status)两张表,且status的默认值为PENDING。 /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 行。/root/outbox/relay.py将status='PENDING'的行 RPUSH 到 Redis 列表outbox.events,并把该行改为PUBLISHED。即使执行两次,列表长度也不应增加。- 同一个
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的大小也必须等于唯一事件数。- 在
/root/outbox/report.txt中写入orders=<n>、outbox=<n>、enqueued=<n>、processed_unique=<n>四行。enqueued必须大于或等于processed_unique。
参考
- 用于查询
PENDING的部分索引:CREATE INDEX ... ON outbox(status) WHERE status='PENDING' - 检查 Redis:
redis-cli LLEN outbox.events、redis-cli HLEN processed - 常见错误 1:试图把中继的发布和状态更新做成一个原子单元——这是不可能的。接受重复并在消费者中将其过滤掉才是正确做法。
- 常见错误 2:不使用
seq,而是按时间戳排序——如果两个事件在同一毫秒内产生,顺序就会颠倒。
创建订单表和发件箱表
使用 /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.py 将 status='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。
将订单数、发件箱行数、进入队列的数量以及实际处理的唯一数量汇总到一个文件中。前三个值与最后一个值之间的关系就是该模式的全部要义。