阻止旧工作进程覆盖新结果:设计原理
一句话总结
在 SQLite 中实现所有权租约和单调递增的 fencing token。
为什么需要它
worker A 领走任务后卡住了。租约过期之后,B 领走了同一个任务并完成了,而后来恢复过来的 A 也写入了完成结果。仅仅规定租约时长,并不能阻止过期 worker 的写入。在写入存储时,也必须确认所有者和代数(generation)。
工作原理
jobs 表保存 id、owner、until、fence、result、done。claim 在 BEGIN IMMEDIATE 内读取当前状态,只领走已经过期的任务,并让 fence 递增。续约和完成,只在当前的 owner 和 fence 都相同、且租约仍然有效时才被允许。它并不保证 exactly-once 执行,但能堵住过去的代数覆盖已保存结果的路径。
claim A fence=1 → 만료 → claim B fence=2 → 완료
A의 fence=1 완료 ───────────→ 거절
阅读契约并预测失败的工作表
下面并不是要求你把实现整个背下来的答案,而是逐步进行的代码评审。每个改动片段都有意破坏了契约。要注意,改动之后正常用例仍然可能通过。执行之前,先预测观测哪些输入、异常、状态能让差异显现出来;实现之后,再拿这个预测与实际结果对比。
1. 创建租约状态表
init_db(path) 以幂等方式创建 jobs(id TEXT PRIMARY KEY, owner TEXT, until REAL NOT NULL DEFAULT 0, fence INTEGER NOT NULL DEFAULT 0, result TEXT, done INTEGER NOT NULL DEFAULT 0)。
判断依据:如果 worker 重启时把代数重置,过期的令牌就会重新变得有效。
需要评审的有问题的改动片段:
CREATE TABLE jobs
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
2. 只登记首次出现的任务
enqueue(path, job_id) 把不存在的 id 以默认状态插入并返回 True;如果已存在,则不改变状态并返回 False。
判断依据:为了让重复登记不会重置正在进行中的租约,要使用 INSERT OR IGNORE。
需要评审的有问题的改动片段:
INSERT OR REPLACE INTO jobs
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
3. 读取当前状态
state(path, job_id) 把行以键为 id、owner、until、fence、result、done 的 dict 返回,不存在时返回 None。
判断依据:通过单独的 DB 连接读取存储的状态,以便与进程内的缓存区分开。
需要评审的有问题的改动片段:
if row else {}
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
4. 以原子方式领走已过期的租约
claim(path, job_id, owner, now, ttl) 在 ttl<=0 时抛出 ValueError。不存在的任务、已完成的任务、until>now 的任务,都返回 None。其余情况保存 owner 和 until=now+ttl,把 fence 加 1,并返回新的 fence。
判断依据:读取和更新要在同一个 BEGIN IMMEDIATE 内完成。边界情况 now==until 时可以重新分配。
需要评审的有问题的改动片段:
row[0] >= now
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
5. 只有当前代数才能续约
renew(path, job_id, owner, fence, now, ttl) 在 ttl<=0 时抛出 ValueError。仅当 owner 和 fence 一致、done=0 且 until>now 时,才把 until 改为 now+ttl 并返回 True,否则返回 False。
判断依据:如果用 renew 把已经过期的所有权救活,就会与新的 worker 冲突。
需要评审的有问题的改动片段:
AND until>=?", (now+ttl
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
6. 拒绝过期的完成
complete(path, job_id, owner, fence, now, result) 仅在租约为当前有效(owner 和 fence 一致、done=0、until>now)时,才保存 result 并设置 done=1,返回 True。其余情况返回 False,并保留原有结果。
判断依据:完成的写入也必须做租约检查,才能挡住迟到的 worker。
需要评审的有问题的改动片段:
AND done>=0 AND until>?", (result
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
7. 只有当前 worker 才能归还租约
release(path, job_id, owner, fence) 对 owner 和 fence 相同且 done=0 的行,把 owner 改为 NULL、until 改为 0,并返回 True。fence 要保留。其余情况返回 False。
判断依据:如果归还时连 fence 也一起重置,就会复用过去的令牌编号。
需要评审的有问题的改动片段:
SET owner=NULL,until=0,fence=0 WHERE
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
8. 并发 claim 时赢家只有一个
claim_many(path, job_id, owners, now, ttl) 用 ThreadPoolExecutor 同时调用每个 owner 的 claim,并按输入顺序返回返回值列表。对于新任务,只有一个 owner 能拿到 fence,其余必须是 None。
判断依据:不要先查询、再关闭连接、然后更新。锁必须同时保护这两个操作。
需要评审的有问题的改动片段:
return [claim(path,job_id,owners[0],now,ttl)]
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
在现场相遇的样子
时钟是由调用方给出的非递减数字,不对多台服务器之间的时钟同步建模。业务的外部副作用并不会被自动 fencing。存储之外的系统同样需要相同的令牌校验,或者单独的幂等处理。SQLite 使用的是真实的锁和事务,但它并不代表大规模分布式队列的吞吐量。
下一项实验要做什么
八个步骤会连成一个可运行的成果。创建租约状态表 → 只登记首次出现的任务 → 读取当前状态 → 以原子方式领走已过期的租约 → 只有当前代数才能续约 → 拒绝过期的完成 → 只有当前 worker 才能归还租约 → 并发 claim 时赢家只有一个。
每个步骤检查的不是函数或文件是否存在,而是实际的返回值、异常和状态变化。看过正确答案之后,故意改动边界比较或清理代码,确认哪些测试会失败。说明为什么前面的测试在后面的步骤中仍然成立,并写出一个本实验不保证的运维条件。