继续完成中断的批次
一句话总结
检查点不是显示在界面上的百分比,而是一个约定:表明已批准列表的哪一部分,已经连同业务变更和审计一起确认。
为什么需要它
外星甜点庆典因暴雨取消。客户要求取消 40 笔未支付订单。前面的课程做的是要么全部成功、要么全部回滚的变更。这次负责人批准说:“可以每 10 笔确认一次,如果后面出问题,请保留已经完成的 10 笔。”仅仅因为这一句话,事务的边界就变了。这不是开发者为了性能而随意拆分的。
长时间一次性加锁,会让其他负责人难以处理订单。反过来,每行都提交,那么即使在分块的第五行出了问题,前四行也无法撤回。合适的大小,不能只由数据库速度来决定。要先问:客户能否理解部分完成?在哪个位置中止,都能给出一致的说明?本实验的 10 笔是教学用的选择,并不是所有生产系统的推荐值。
工作原理
1. 把批准列表和执行结果分开
注册时保存 job_id、tenant、每笔订单的 id、revision、qty 以及 chunk_size。批准列表按 ID 规范化排序,执行过程中不改变。用同一个作业 ID,只改变列表顺序重新注册,仍是同一个请求。给同一个 ID 附上不同的客户、数量、分块大小,就是冲突。悄悄覆盖的话,以前的进度就会指向新的列表。
如果在执行过程中重新检索当前 pending 订单来构造列表,会怎样?第一个分块结束之后若有新订单进来,最初并不存在的订单就可能混入取消对象。订单被删除的话,OFFSET 的含义也会变。所以这里使用的不是 WHERE 结果的第几行,而是冻结的批准数组的 next_index。即使数据从 40 笔变成 39 笔,也不会把批准本身改写成 39 笔。
2. 锁住作业行,挑选下一个分块
两个 worker 同时读到 next_index=10,就都可能拿到第二个分块。用 SELECT FOR UPDATE 读取 jobs 中对应的行,并持有到事务结束,就能把同一作业的下一个区间选择串行化。PostgreSQL 的行锁在事务结束时释放,并不是阻止所有一般查询的全局锁。请看官方行锁文档。
在这个设计里,两个 worker 并不会同时处理同一作业的不同分块。两个 worker 是为了确认故障之后的接续,以及防止重复选择。如果需要更高的并行吞吐量,就需要按分块的租约、过期、所有权代这类另外的设计。不要说当前代码有这样的保证。
3. 把三条记录捆绑在一个分块里
如果分块是第 11–20 号订单,下面的变化是同一个事务。
| 保存对象 | 要确认的内容 | 分开提交会出现的问题 |
|---|---|---|
| orders | 把批准版本的 pending 改为 cancelled,并增加版本 | 只有订单变了,恢复位置还是原样 |
| job_audit | 序号、ID、前后版本、数量 | 只有审计,而实际订单没有变 |
| jobs.next_index | 下一个区间的起点 20 | 把没有执行的订单误认为已完成 |
UPDATE 条件里要放入客户、ID、版本、数量、状态全部。RETURNING 给出实际被修改的行的新值。没有匹配的行时,SQL 本身并不报错,所以程序要把它变成 Conflict,让整个分块回滚。请把 UPDATE 官方说明中的返回值和受影响行数一起读。
Python 的低层函数也使用事务,但在外层 run_chunk 内被调用时,不能提前提交整体。psycopg 嵌套的 transaction 上下文是用 SAVEPOINT 工作的。本课采用的方式,是在 autocommit=True 的连接上显式开启外层事务边界。请对照 psycopg 事务说明和你自己的代码,确认最外层的上下文在哪里。
4. 恢复之前要怀疑进度
不能因为 next_index 是 20,就全盘相信。审计的序号以及 ID、前后版本、数量,必须与批准数组的前 20 个完全一致。仅凭有 20 行这一事实,无法知道有没有夹进别的 ID。本实验拒绝没有审计的进度、指向分块中间的进度、超过列表长度的进度。不会自动“修补”损坏的记录而去修改更多订单。
第一个分块确认之后,如果第二个分块里的 15 号订单被其他负责人修改,就要把第二个分块里 11–14 号的变更一起回滚。第一个分块的记录和 15 号的后续变更保留。读取新的 revision、当场重新做出批准,不是重试,而是新的业务决定。必须中止,说明剩余范围之后,重新获得批准。
在现场相遇的样子
丢失了响应后再次调用,会返回什么
假设分块提交成功了,但客户端在收到结果之前终止了。恢复同一个作业时,不会重放刚完成的分块的答案,而是处理下一个分块。因为这个 API 的契约不是“执行第 N 个分块”,而是“推进这个作业的下一个未完成分块”。所以不能只把返回的 processed 加起来,作为整个作业完成的证据。完整的历史要从数据库的审计和检查点中读取。
max_chunks 是一次 drain 调用所尝试的分块调用次数。它不等同于每次调用的 SQL 时间上限或总耗时。如果其他 worker 已经完成了大部分,我的调用也可能收到处理 ID 为空的完成响应。正常完成就立即停止,对冲突或连接错误,也不是无条件重试。不要忘了前面课程学到的错误分类。
向客户把“完成”和“当前状态”区分开来说明
过去完成的订单,即使之后又被修改或删除,当时完成的事实也不会消失。报告要把批准 ID、已确认 ID、未处理 ID 分开,并把已确认 ID 的当前状态分为 matching、drifted、missing。不能因为找到了数量相同的另一笔订单,就把 missing 去掉。把进度和当前行用一条 SELECT 读取,就能在 Read Committed 的语句级快照下进行核对。请参考隔离级别官方文档。
FDE 要把客户的业务条件转换成实现的不变式,并且在失败时说明确认到了哪里。所查看的 Palantir FDE 招聘启事强调理解客户问题并实现真实的解决方案。这个虚构案例是作者为练习这种能力而设计的,并不意味着该启事必须要求 PostgreSQL 或这种实现。
下一项实验要做什么
从 40 笔批准、10 笔分块开始,让它在其他客户、不连续的 ID、最后一个较短的分块上也能工作。在订单变更之后、审计之后、检查点之后、提交之后这四个点真实终止客户端,并用新连接确认记录。让两个进程接续同一个作业,并用独立的 SQL 核对精确的批准 ID 是否都只被处理过一次。
前提也要明确。正常的订单写入会让 revision 增加,已确认的批准和审计不可变。至于连实验数据库账户的恶意直接修改也要防止的权限设计,是另外的事。服务器保持运行,只终止客户端,所以并没有验证服务器断电的耐久性、外部支付退款、无限制的吞吐量。