两名零食配送员与无休止的重试
目标
即使两名节日零食投递方相互竞争,且其中一名在发送之后立即退出,也要在有限的预算内恢复剩余的业务。要让过期的工作者无法覆盖新工作者的完成记录。
为什么重要
网络超时并不意味着业务没有被处理。即使租约已经结束,之前的请求仍可能继续执行。这次要把队列的先占令牌与接收端 inbox 区分开,并把锁之外的发送、失败预算和隔离记录一起设计。把前面 outbox 实验中的重复效果防护,扩展到多个投递方的场景。
预计需要 110 分钟。比默认会话更长,所以请在到期之前点击“+时间”按钮延长。会话结束后文件会消失,需要的代码请另行保存。开始前应已了解 Python 异常处理、SQLite 事务,以及前面模块中的 outbox 和接收端去重。
数据契约
交付物是 /root/lease/worker.py。队列文件和 HTTP 接收服务器由检查器临时创建并清理,所以不要把路径和端口写死。所有队列都是具有正常 schema 的可信本地文件。使用下面的 schema。
CREATE TABLE jobs (seq INTEGER PRIMARY KEY AUTOINCREMENT,
id TEXT NOT NULL UNIQUE, qty INTEGER NOT NULL,
status TEXT NOT NULL CHECK(status IN ('pending','leased','done','dead')),
attempts INTEGER NOT NULL, max_attempts INTEGER NOT NULL, token INTEGER NOT NULL,
available_at INTEGER NOT NULL, deadline INTEGER NOT NULL,
owner TEXT, lease_until INTEGER, reason TEXT);
CREATE TABLE redrives (seq INTEGER PRIMARY KEY AUTOINCREMENT,
id TEXT NOT NULL, token INTEGER NOT NULL, at_ms INTEGER NOT NULL, note TEXT NOT NULL);
标识符是由 ASCII 字母、数字、下划线、连字符组成、长度为 1–64 的严格 str。now_ms 和 clock 的返回值是以统一时间基准表示的严格 int,范围从 0 到 (10**15-60000),并且所有整数契约都拒绝 bool。各个 DB 函数会在修改之前以 ValueError 拒绝类型和范围错误,并且不修改输入。如果 run_once 的第二次 clock 无效或时间倒退,则保留已经提交的先占并传递错误。这并不是在主机之间校准时钟的共识算法。
借来的 con 在调用开始时没有事务,函数也不会关闭它。无论成功还是失败,都不留下未结束的事务。写入时的选择与修改用 BEGIN IMMEDIATE 绑定,在提交之前出现 fault 错误时整体回滚并传递原始错误。提交之后的错误则保留已经确定的状态。fault=None 时不调用,钩子只在执行了相应修改的路径上调用。在返回 None、False 等没有修改的路径上没有钩子。即使没有目标,claim 也会提交前面已执行的隔离更新。
每个任务的含义是相互独立的。这不是前一个业务失败就会阻塞后一个业务的严格顺序保证队列。token 是按业务单调递增的先占编号,检查范围内不超过 10**15。队列的标记检查并不能替代外部服务器的权限验证或对外部资源的 fencing。提供的接收服务器会去除相同 ID、相同数量的重复效果。redrive 的前提是由有权限的运维人员调用,没有实现登录和角色管理功能。
步骤
- 计算重试延迟的上限与抖动——在 worker.py 中定义 Exception 的子类 Conflict、Retryable、Permanent、BadAck,并实现 retry_delay(attempt,base_ms,cap_ms,jitter)。attempt 是 1–16,base_ms 是 1–60000,cap_ms 是 base_ms–60000,均为严格的 int。jitter 是不含 bool 的严格 int/float,且是有限的 0–1。返回 floor(min(cap_ms,base_ms*2(attempt-1))*jitter)。无效输入为 ValueError。不要使用内部随机数或 sleep,按给定的比例计算。
- 打开重启后依然保留的队列——open_queue(path) 仅在不存在时,用一个事务创建下面的 jobs、redrives 表,并返回 sqlite3.Connection。使用 isolation_level=None、timeout=1 秒、journal_mode=DELETE、synchronous=FULL。保留已有的业务行和审计行,初始化失败时关闭连接。接收的是一次性的本地文件,而不是内存 DB。
- 不因重复受理而重置预算——enqueue(con,event,now_ms,ttl_ms,max_attempts=3) 接收只含 id、qty 的严格 dict。id 是下面的标识符,qty 是严格的 int 1–1000,ttl_ms 是 1–60000,max_attempts 是 1–16。把新业务以 pending、attempts=0、token=0、available_at=now_ms、deadline=now_ms+ttl_ms、owner、lease_until、reason=NULL 插入,并返回 True。相同 ID、相同数量,即使输入的预算不同,也会保留已有的整行并返回 False;数量不同则是 Conflict。先验证所有输入,并把选择和插入放在同一个写事务中处理。
- 不让两个投递方同时领走同一个任务——实现 claim(con,owner,now_ms,lease_ms=1000,fault=None)。owner 是标识符,lease_ms 是严格的 int 1–60000。在 BEGIN IMMEDIATE 内,把 pending 或已过期的 leased 行中满足 deadline<=now_ms 或 attempts>=max_attempts 的行移到 dead。reason 在截止时为 deadline,否则为 exhausted,并把 owner、lease_until 设为 NULL。在剩下的行中,按 seq 升序选出一个业务:状态为 pending 且 available_at<=now_ms,或状态为 leased 且 lease_until<=now_ms。没有则返回 None。有的话,把 attempts、token 各加 1,并保存为 leased、owner、lease_until=min(now_ms+lease_ms,deadline)、reason=NULL。返回只包含 id、qty、owner、token、attempt、lease_until 的 dict。attempt 是更新后的 attempts。更新之后调用 fault('after-claim'),COMMIT 之后调用 fault('after-commit')。
- 拒绝过期投递方的完成——实现 finish(con,ticket,outcome,now_ms,delay_ms=0,fault=None)。ticket 是只含 claim 那 6 个键的 dict,其中 id、owner 是标识符,qty=1–1000,token=1–1015,attempt=1–16,lease_until=1–10**15,均为严格的 int。outcome 是 ok、retry、permanent 之一的 str,delay_ms 是严格的 int 0–60000。在一个写事务中,如果业务不存在、不是 leased、标记的数量、所有者、令牌、尝试次数、到期时间不一致,或 now_ms>=lease_until 或 deadline,则不做任何修改并返回 False。有效的 ok 为 done 且原因为 NULL,permanent 为 dead/permanent。retry 在已达到最大尝试次数时为 dead/exhausted,否则若 now_ms+delay_ms>=deadline 则为 dead/deadline,其余为 pending/retry。仅当为 pending 时 available_at=now_ms+delay_ms,其余情况为 now_ms,并把 owner、lease_until 设为 NULL。保留其他字段,并返回新的 status 字符串。更新之后调用 fault('after-finish'),COMMIT 之后调用 fault('after-commit')。
- 同时留下解除隔离与审计记录——redrive(con,job_id,now_ms,ttl_ms,note,fault=None) 只对存在的 dead 业务重新驱动,并返回 True。job_id、note 是标识符,ttl_ms 是严格的 int 1–60000。目标不存在或不是 dead 时为 ValueError。在一个写事务中,顺序为:向 redrives 插入 id、当前 token、at_ms=now_ms、note → fault('after-audit') → 把业务改为 pending、attempts=0、available_at=now_ms、deadline=now_ms+ttl_ms、owner、lease_until、reason=NULL → fault('after-redrive') → COMMIT → fault('after-commit')。保留 qty、max_attempts、token。
- 在写锁之外执行网络发送——run_once(con,owner,clock,send,jitter,lease_ms=1000) 先验证 jitter,并用 clock() 的第一个时刻执行 claim。没有业务则返回 None,且不调用 send。有业务时,在事务之外把只含 id、qty 的新 dict 向 send 传递一次。只有严格的 True 才算 ok,其他返回值为 BadAck。只把 Retryable 归类为 retry,并使用 retry_delay(ticket.attempt,100,1000,jitter);Permanent 归类为 permanent。其他错误和 BadAck 原样传递,并保留 leased 状态。在正常结束或已分类的错误之后再次调用 clock(),如果第二个值小于第一个值则为 ValueError。把第二个时刻以及结果、延迟传给 finish 之后,返回 id、token、status 的 dict。如果 finish 返回 False,则 status 为 stale。即使发送回调修改了收到的 dict,或通过单独的连接放入了业务,当前标记也必须保持不变。
- 由另一个进程恢复发送之后就死掉的任务——run_file(path,owner,clock,send,jitter,lease_ms=1000) 在已有的普通队列文件不存在时,以 FileNotFoundError 拒绝,不会新建空队列。它通过 open_queue 拥有连接并返回 run_once 的结果,无论成功还是失败都关闭连接。最终检查会验证:真实的两个子进程同时 claim 时只有一个成功;在提交前后被强制终止之后,状态是原子的。单独的 HTTP 接收服务器在提交数量之后立即终止第一个投递方,并由第二个投递方进程重新发送。HTTP 请求有两次,但 inbox 和接收效果必须只有一次。接收服务器和临时 DB 由检查器提供。
参考
- 直接诊断:python3 -B /opt/fixtures/lease/check.py 8 /root/lease/worker.py。把 8 改成当前步骤,就会检查到该步骤。每次评分上限为 12 秒,不会修改提交的文件。
- 发送回调是带有自身超时的适配器。run_once 不会在内部无限重试或 sleep。租约到期不会取消正在执行的请求。
- 只使用标准库、小型本地 DB 和回环 HTTP。不需要互联网安装、外部服务或额外的 capability。会测试真实的子进程终止,但不保证断电、磁盘丢失、恶意替换文件以及服务器之间的时钟共识。
- SQLite rollback journal 会把简短的先占写入串行化。这与前面快照模块中“读取期间的写入时点隔离”要求不同,是另一种工作负载。不要把 DB 锁延长到网络发送的时间。
计算重试延迟的上限与抖动
在 worker.py 中定义 Exception 的子类 Conflict、Retryable、Permanent、BadAck,并实现 retry_delay(attempt,base_ms,cap_ms,jitter)。attempt 是 1–16,base_ms 是 1–60000,cap_ms 是 base_ms–60000,均为严格的 int。jitter 是不含 bool 的严格 int/float,且是有限的 0–1。返回 floor(min(cap_ms,base_ms*2**(attempt-1))*jitter)。无效输入为 ValueError。不要使用内部随机数或 sleep,按给定的比例计算。
bool 在 Python 中是 int 的子类型。先应用上限,再乘以注入的比例,然后向下取整。
打开重启后依然保留的队列
open_queue(path) 仅在不存在时,用一个事务创建下面的 jobs、redrives 表,并返回 sqlite3.Connection。使用 isolation_level=None、timeout=1 秒、journal_mode=DELETE、synchronous=FULL。保留已有的业务行和审计行,初始化失败时关闭连接。接收的是一次性的本地文件,而不是内存 DB。
只用简短的写事务来协调工作者。不要把创建表和初始化已有内容混为一谈。
不因重复受理而重置预算
enqueue(con,event,now_ms,ttl_ms,max_attempts=3) 接收只含 id、qty 的严格 dict。id 是下面的标识符,qty 是严格的 int 1–1000,ttl_ms 是 1–60000,max_attempts 是 1–16。把新业务以 pending、attempts=0、token=0、available_at=now_ms、deadline=now_ms+ttl_ms、owner、lease_until、reason=NULL 插入,并返回 True。相同 ID、相同数量,即使输入的预算不同,也会保留已有的整行并返回 False;数量不同则是 Conflict。先验证所有输入,并把选择和插入放在同一个写事务中处理。
受理重试不是新业务。请比较 UNIQUE ID 与已有数量,也不要修改输入的 dict。
不让两个投递方同时领走同一个任务
实现 claim(con,owner,now_ms,lease_ms=1000,fault=None)。owner 是标识符,lease_ms 是严格的 int 1–60000。在 BEGIN IMMEDIATE 内,把 pending 或已过期的 leased 行中满足 deadline<=now_ms 或 attempts>=max_attempts 的行移到 dead。reason 在截止时为 deadline,否则为 exhausted,并把 owner、lease_until 设为 NULL。在剩下的行中,按 seq 升序选出一个业务:状态为 pending 且 available_at<=now_ms,或状态为 leased 且 lease_until<=now_ms。没有则返回 None。有的话,把 attempts、token 各加 1,并保存为 leased、owner、lease_until=min(now_ms+lease_ms,deadline)、reason=NULL。返回只包含 id、qty、owner、token、attempt、lease_until 的 dict。attempt 是更新后的 attempts。更新之后调用 fault('after-claim'),COMMIT 之后调用 fault('after-commit')。
不要在选择与更新之间释放锁。要区分到期边界上的等号,以及即使没有可选的任务也必须确定的隔离变更。
拒绝过期投递方的完成
实现 finish(con,ticket,outcome,now_ms,delay_ms=0,fault=None)。ticket 是只含 claim 那 6 个键的 dict,其中 id、owner 是标识符,qty=1–1000,token=1–1015,attempt=1–16,lease_until=1–1015,均为严格的 int。outcome 是 ok、retry、permanent 之一的 str,delay_ms 是严格的 int 0–60000。在一个写事务中,如果业务不存在、不是 leased、标记的数量、所有者、令牌、尝试次数、到期时间不一致,或 now_ms>=lease_until 或 deadline,则不做任何修改并返回 False。有效的 ok 为 done 且原因为 NULL,permanent 为 dead/permanent。retry 在已达到最大尝试次数时为 dead/exhausted,否则若 now_ms+delay_ms>=deadline 则为 dead/deadline,其余为 pending/retry。仅当为 pending 时 available_at=now_ms+delay_ms,其余情况为 now_ms,并把 owner、lease_until 设为 NULL。保留其他字段,并返回新的 status 字符串。更新之后调用 fault('after-finish'),COMMIT 之后调用 fault('after-commit')。
仅凭 owner 相同,也可能是过期的运行。要检查整个标记和时间边界,不要因为结果来得晚就删除当前业务。
同时留下解除隔离与审计记录
redrive(con,job_id,now_ms,ttl_ms,note,fault=None) 只对存在的 dead 业务重新驱动,并返回 True。job_id、note 是标识符,ttl_ms 是严格的 int 1–60000。目标不存在或不是 dead 时为 ValueError。在一个写事务中,顺序为:向 redrives 插入 id、当前 token、at_ms=now_ms、note → fault('after-audit') → 把业务改为 pending、attempts=0、available_at=now_ms、deadline=now_ms+ttl_ms、owner、lease_until、reason=NULL → fault('after-redrive') → COMMIT → fault('after-commit')。保留 qty、max_attempts、token。
尝试预算重新给,但先占代次不能复用。要防止只留下审计行或只解除状态这类做了一半的重新驱动。
在写锁之外执行网络发送
run_once(con,owner,clock,send,jitter,lease_ms=1000) 先验证 jitter,并用 clock() 的第一个时刻执行 claim。没有业务则返回 None,且不调用 send。有业务时,在事务之外把只含 id、qty 的新 dict 向 send 传递一次。只有严格的 True 才算 ok,其他返回值为 BadAck。只把 Retryable 归类为 retry,并使用 retry_delay(ticket.attempt,100,1000,jitter);Permanent 归类为 permanent。其他错误和 BadAck 原样传递,并保留 leased 状态。在正常结束或已分类的错误之后再次调用 clock(),如果第二个值小于第一个值则为 ValueError。把第二个时刻以及结果、延迟传给 finish 之后,返回 id、token、status 的 dict。如果 finish 返回 False,则 status 为 stale。即使发送回调修改了收到的 dict,或通过单独的连接放入了业务,当前标记也必须保持不变。
不要在等待外部响应时占着 DB 事务。要用第二个时刻重新判断有效性,不要把未知的失败掩盖成成功。
由另一个进程恢复发送之后就死掉的任务
run_file(path,owner,clock,send,jitter,lease_ms=1000) 在已有的普通队列文件不存在时,以 FileNotFoundError 拒绝,不会新建空队列。它通过 open_queue 拥有连接并返回 run_once 的结果,无论成功还是失败都关闭连接。最终检查会验证:真实的两个子进程同时 claim 时只有一个成功;在提交前后被强制终止之后,状态是原子的。单独的 HTTP 接收服务器在提交数量之后立即终止第一个投递方,并由第二个投递方进程重新发送。HTTP 请求有两次,但 inbox 和接收效果必须只有一次。接收服务器和临时 DB 由检查器提供。
下一个进程会重新打开同一个文件。不要声称已经消除了外部接收效果与队列完成之间的缝隙,而要恢复重复发送。