TT Lab
开始
学习 学习路径 课程

不可撤销的变更

外星人节庆的取消任务停在了25%

在 TT Lab 中继续学习

目标

把外星甜点庆典的 40 笔订单,每 10 笔一组地取消。批准列表固定,即使中途进程消失,也从已确认的分块之后继续。

为什么重要

把整体回滚换成逐分块确认,业务契约也会随之改变。必须明确地获得对部分完成的批准,并检查检查点是否与实际的业务变更和审计相符。请先学习 Python 函数和异常、SQL 事务以及前面的批准版本、补偿实验。预计 120 分钟,所以请在到期前用“+时间”延长。会话结束后文件会消失。需要的代码请另外保存。

环境与通用契约

产出物是 /root/chunks/worker.py。镜像里有 PostgreSQL 16、psycopg 3.2.3、Python 3,没有运行时安装。可以用 postgres 用户写入 /root,不需要额外的 capability 或用户切换。

评分器会在本地 labdb 的单独临时 schema 中准备下面的表和虚构订单,并且只清理自己建立的 schema。学员函数使用传入连接的 search_path 和 DSN。不要修改 public 表,也不要把 schema、客户、ID、DSN 写死。SQL 的值通过参数传递。

CREATE TABLE orders(id integer PRIMARY KEY,tenant text NOT NULL,
 qty integer NOT NULL CHECK(qty BETWEEN 1 AND 1000),
 state text NOT NULL CHECK(state IN ('pending','paid','cancelled')),
 revision integer NOT NULL CHECK(revision>=0));
CREATE TABLE jobs(job_id text PRIMARY KEY,tenant text NOT NULL,targets jsonb NOT NULL,
 chunk_size integer NOT NULL CHECK(chunk_size BETWEEN 1 AND 10),
 next_index integer NOT NULL CHECK(next_index>=0));
CREATE TABLE job_audit(job_id text REFERENCES jobs(job_id),ordinal integer NOT NULL,
 id integer NOT NULL,previous_revision integer NOT NULL,new_revision integer NOT NULL,
 qty integer NOT NULL,PRIMARY KEY(job_id,ordinal),UNIQUE(job_id,id));

标识符 job_id 和 tenant 是 str 本身(类型必须恰好是 str),是 ASCII 英文、数字、下划线、连字符 1–64 个字符。items/targets 是 1–64 个元素的 list 本身,每个元素是只有 id、revision、qty 的 dict 本身。id 是 1–2147483647 的 int,revision 是 0–2147483646 的 int,qty 是 1–1000 的 int。bool 不允许当作整数。拒绝重复的 ID,并规范化为按 ID 排序的深拷贝。chunk_size 是 1–10 的 int(类型必须恰好是 int),allow_partial 必须恰好是 True。错误的直接输入,在写入之前以 ValueError 报错,不自动修改。

契约是:注册之后,jobs 的 tenant、targets、chunk_size 和 job_audit 不可变。next_index 在 0 以上、不超过对象数;如果不是完成位置,就必须是 chunk_size 的倍数。保存的 targets 必须与规范化输入相同。审计按 ordinal 顺序,必须与原批准数组的 [0:next_index] 在 ID、之前 revision、新 revision=之前+1、qty 上完全相同。多出来的审计也是损坏。原订单的后续变更不是过去审计的损坏,所以不要用前面分块的当前值去覆盖过去的批准。

run_chunk 要处理的区间是 [next_index:min(next_index+chunk_size,对象数)]。fault 只有存在时才用各个字符串参数调用。不要隐藏钩子错误。锁和语句上限以每条 SQL 为准,不是整个分块耗时的限制。成功和失败之后,请恢复原连接的设置。

借来的 con 是 autocommit=True、Read Committed,外部调用开始时没有打开的事务。函数不关闭连接,成功和失败之后都不留下事务。cancel_chunk 内部的嵌套调用,必须保留外层事务。run_chunk 会在真正提交之后调用钩子,所以不要把它包在其他外部事务里。只有 chunk_file 拥有连接。

步骤

  1. 固定部分完成的批准——实现继承 Exception 的 Conflict 和 manifest(job_id,tenant,items,chunk_size,allow_partial)。校验下面的输入契约,并返回 job_id、tenant、targets、chunk_size、allow_partial 的新 dict。targets 是按 ID 排序的新 list 和新 dict。只允许明确的 True 批准,错误的输入是 ValueError。
  2. 重新注册也不覆盖批准——register(con,plan) 校验并规范化只含 manifest 的五个键的 dict 本身之后,以 next_index=0 注册到 jobs,并返回 True。同一个作业 ID 的相同客户、规范化 targets、chunk_size,什么也不改变并返回 False,内容不同则是 Conflict。输入错误是写入之前的 ValueError,不修改订单和审计。
  3. 分别读取批准和进度——inspect_job(con,job_id) 校验 ID,没有则返回 None,有则返回 job_id、tenant、targets、chunk_size、next_index 的 dict。只读,修改返回值不影响原记录。
  4. 分块内一行失败也要全部回滚——cancel_chunk(con,tenant,items) 先校验对象输入,按 ID 顺序,用当前的客户、ID、revision、qty、pending 作条件做条件 UPDATE。改为 cancelled,把 revision 加 1,返回 id、previous_revision、new_revision、qty 的 dict 列表。只要有一行不一致,就以 Conflict 回滚整个分块。不写 jobs 和 job_audit,嵌套调用不提前确认外层事务。
  5. 把变更、审计、位置一起确认——run_chunk(con,job_id,fault=None) 先设置自己事务的锁上限 500ms、语句上限 2000ms,再锁住 jobs 行来读取。如果不存在,或者保存的批准、next_index、审计与下面的契约不同,就是 Conflict。用 cancel_chunk 变更下一个区间之后调用 fault(after-orders),把该区间的序号、ID、前后版本、数量全部保存到 job_audit 之后调用 fault(after-audit),把 next_index 改为区间末尾之后调用 fault(after-checkpoint),真正 COMMIT 之后调用 fault(after-commit)。返回值是按 ID 排序的 processed list、next_index、done bool。如果已经完成,则 processed=[]、done=True,并且不调用钩子。提交之前出错,只回滚本次分块并保留前面的分块;提交之后出错,保留已确认的状态并传递原始错误。
  6. 报告过去的完成和当前的差异——report(con,job_id) 没有则是 Conflict,有则返回含 approved、committed、remaining、matching、drifted、missing 的按 ID 排序的 list 的 dict。approved 是原批准,committed 是 next_index 之前的部分,remaining 是其余部分。committed 的当前行在客户、qty、cancelled、原 revision+1 上全部相同则是 matching,ID 不存在是 missing,其余是 drifted。进度和当前行用一条 SELECT 查询,不修改数据。
  7. 客户端终止之后接续剩余区间——chunk_file(dsn,job_id,fault=None) 用 psycopg.connect(dsn,autocommit=True,connect_timeout=2) 打开自己拥有的连接,调用 run_chunk 并返回同样的结果。无论成功还是失败都关闭连接,不隐藏错误。检查四个钩子点真实终止客户端之后的恢复结果。
  8. 用两个 worker 只推进限定的次数——drain(dsn,job_id,max_chunks) 在建立连接之前,校验 ID 和类型恰好为 int 的 1–64 的 max_chunks。最多调用 max_chunks 次 chunk_file,但 done=True 就立即结束。返回调用次数 calls、本次调用实际处理的 processed 的合并列表、最后的 next_index、最后的 done。任何错误都立即传递,不做隐藏的重试。两个进程把同一个作业处理到底时,全部批准 ID 必须恰好各被修改一次。

参考

固定部分完成的批准

实现继承 Exception 的 Conflict 和 manifest(job_id,tenant,items,chunk_size,allow_partial)。校验下面的输入契约,并返回 job_id、tenant、targets、chunk_size、allow_partial 的新 dict。targets 是按 ID 排序的新 list 和新 dict。只允许明确的 True 批准,错误的输入是 ValueError。

“40 笔全部原子”的约定,与“每 10 笔确认”的约定是不同的。bool 和 int 也要区分。

重新注册也不覆盖批准

register(con,plan) 校验并规范化只含 manifest 的五个键的 dict 本身之后,以 next_index=0 注册到 jobs,并返回 True。同一个作业 ID 的相同客户、规范化 targets、chunk_size,什么也不改变并返回 False,内容不同则是 Conflict。输入错误是写入之前的 ValueError,不修改订单和审计。

把已有的 next_index 覆盖为 0 的 upsert,不是重新注册,而是进度丢失。

分别读取批准和进度

inspect_job(con,job_id) 校验 ID,没有则返回 None,有则返回 job_id、tenant、targets、chunk_size、next_index 的 dict。只读,修改返回值不影响原记录。

不要把批准数组和当前订单混为一谈。尚未执行的注册,也是有效的记录。

分块内一行失败也要全部回滚

cancel_chunk(con,tenant,items) 先校验对象输入,按 ID 顺序,用当前的客户、ID、revision、qty、pending 作条件做条件 UPDATE。改为 cancelled,把 revision 加 1,返回 id、previous_revision、new_revision、qty 的 dict 列表。只要有一行不一致,就以 Conflict 回滚整个分块。不写 jobs 和 job_audit,嵌套调用不提前确认外层事务。

UPDATE 受影响的行数为 0,SQL 也是成功的。请判断没有返回行是不是业务冲突。

把变更、审计、位置一起确认

run_chunk(con,job_id,fault=None) 先设置自己事务的锁上限 500ms、语句上限 2000ms,再锁住 jobs 行来读取。如果不存在,或者保存的批准、next_index、审计与下面的契约不同,就是 Conflict。用 cancel_chunk 变更下一个区间之后调用 fault(after-orders),把该区间的序号、ID、前后版本、数量全部保存到 job_audit 之后调用 fault(after-audit),把 next_index 改为区间末尾之后调用 fault(after-checkpoint),真正 COMMIT 之后调用 fault(after-commit)。返回值是按 ID 排序的 processed list、next_index、done bool。如果已经完成,则 processed=[]、done=True,并且不调用钩子。提交之前出错,只回滚本次分块并保留前面的分块;提交之后出错,保留已确认的状态并传递原始错误。

请把审计序号与批准数组的前缀集合比较。同一个作业行的锁,从选择下一个区间之前就需要。

报告过去的完成和当前的差异

report(con,job_id) 没有则是 Conflict,有则返回含 approved、committed、remaining、matching、drifted、missing 的按 ID 排序的 list 的 dict。approved 是原批准,committed 是 next_index 之前的部分,remaining 是其余部分。committed 的当前行在客户、qty、cancelled、原 revision+1 上全部相同则是 matching,ID 不存在是 missing,其余是 drifted。进度和当前行用一条 SELECT 查询,不修改数据。

当前其他 ID 的相同值,无法代替缺失的原 ID。推进下一个分块,不能覆盖前面分块的后续变更。

客户端终止之后接续剩余区间

chunk_file(dsn,job_id,fault=None) 用 psycopg.connect(dsn,autocommit=True,connect_timeout=2) 打开自己拥有的连接,调用 run_chunk 并返回同样的结果。无论成功还是失败都关闭连接,不隐藏错误。检查四个钩子点真实终止客户端之后的恢复结果。

如果提交之后丢失了响应,同一作业的下一次推进,不是已经结束的分块,而是下一个分块。

用两个 worker 只推进限定的次数

drain(dsn,job_id,max_chunks) 在建立连接之前,校验 ID 和类型恰好为 int 的 1–64 的 max_chunks。最多调用 max_chunks 次 chunk_file,但 done=True 就立即结束。返回调用次数 calls、本次调用实际处理的 processed 的合并列表、最后的 next_index、最后的 done。任何错误都立即传递,不做隐藏的重试。两个进程把同一个作业处理到底时,全部批准 ID 必须恰好各被修改一次。

把不同 worker 返回的 ID,与独立的数据库审计比较。为了拿到完成响应而做的空调用,也算在 calls 里。