从最后提交的检查点恢复导入:设计原理
一句话总结
把源文件指纹与检查点结合起来,防止错误地接着处理另一个文件。
为什么需要它
一个正在导入几千行数据的任务中途崩溃了。运维人员换了文件,用同一个任务 id 重新运行,结果前半部分来自旧文件,后半部分来自新文件。只保存已处理的行号,就无法确认输入的身份。必须把源文件字节的指纹和最后提交的位置一起保存。
工作原理
输入是没有重复 id 的 JSON 数组,每一行都有 id 和 value。把源 bytes 的 SHA-256 和检查点保存在 imports 中。批次中每个 item 的插入与 next_index 的递增,在同一个事务中完成。如果某一行中途抛出异常,整个批次都会回滚,而之前已经完成的批次保留下来。重新执行时从保存的索引开始,但如果源文件的指纹不同,就拒绝。
원본 bytes → 지문 확인 → next_index → 배치 INSERT + 체크포인트 COMMIT
└ 중간 실패 → 이번 배치만 ROLLBACK
阅读契约并预测失败的工作表
下面并不是要求你把实现整个背下来的答案,而是逐步进行的代码评审。每个改动片段都有意破坏了契约。要注意,改动之后正常用例仍然可能通过。执行之前,先预测观测哪些输入、异常、状态能让差异显现出来;实现之后,再拿这个预测与实际结果对比。
1. 确认输入行的契约
parse_rows(raw) 读取 JSON 数组的 bytes。每一项都要有 id(非空 str)和 value(不含 bool 的 int),并且禁止 id 重复。违反时抛出 ValueError。返回只包含 {id,value} 的行列表。
判断依据:如果悄悄地用最后一行覆盖重复的 id,导入结果就无法预测。
需要评审的有问题的改动片段:
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
2. 固定源文件字节的指纹
source_digest(raw) 是对 bytes 应用 SHA-256 得到的 hex 字符串。不对 JSON 做规范化。
判断依据:这份契约只允许针对同一个源文件来恢复执行。
需要评审的有问题的改动片段:
hashlib.sha256(raw.replace(b" ",b""))
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
3. 保存输入和检查点
init_db(path) 以幂等方式创建 imports(id TEXT PRIMARY KEY,digest TEXT NOT NULL,next_index INTEGER NOT NULL) 和 items(batch TEXT NOT NULL,id TEXT NOT NULL,value INTEGER NOT NULL,PRIMARY KEY(batch,id))。
判断依据:不同导入任务中相同的行 id 要分开保存。
需要评审的有问题的改动片段:
CREATE TABLE imports
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
4. 防止用另一个源文件恢复执行
begin(path,batch,digest) 对新任务会保存 next_index=0 并返回 0;已有的任务如果指纹相同,就返回 next_index;指纹不同则抛出 ValueError。
判断依据:即使任务 id 和行号相同,输入文件也可能不同。
需要评审的有问题的改动片段:
if False:
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
5. 以原子方式提交批次和位置
apply_chunk(path,batch,start,rows,fault=lambda index:None) 只有在当前 next_index==start 时才执行,否则抛出 ValueError。把 rows 按顺序放入 items,每次插入之后调用 fault(全局索引)。全部成功后,保存 next_index=start+len(rows) 并返回。
判断依据:如果每一行都单独提交,检查点与行的状态就会不一致。
需要评审的有问题的改动片段:
db.commit()
fault(start+offset)
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
6. 查询当前位置
checkpoint(path,batch) 返回 next_index,任务不存在时返回 None。
判断依据:读取的是最后一次提交的位置,而不是最后一次尝试处理的位置。
需要评审的有问题的改动片段:
return 0 if row else None
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
7. 按任务分开结果
values(path,batch) 按 id 升序返回该 batch 的 (id,value) 元组。
判断依据:把 batch 设为条件,避免其他任务中相同 id 的行混进结果。
需要评审的有问题的改动片段:
WHERE batch!=? ORDER BY id
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
8. 中途失败之后安全地接着处理
import_all(path,batch,raw,size=2,fault=lambda index:None) 校验 size 是(不含 bool 的)正 int,并使用 parse_rows、source_digest 和 begin。把剩余的行按每 size 行一批交给 apply_chunk 处理,并返回最终的 checkpoint。
判断依据:要确认这样的场景:保留第一个批次的成功结果,在第二个批次失败之后恢复执行。
需要评审的有问题的改动片段:
begin(path,batch,source_digest(raw))
index=0
将它与包含该片段的函数的公开契约对照。如果仅凭一个成功用例无法区分,就选择本应被拒绝的输入,或失败之后的状态作为观测对象。
在现场相遇的样子
这是面向小规模数据的实验,会把整个输入读进内存。不要把它夸大成能够流式解析大文件的引擎。这是一份保守的契约:哪怕源文件只有空白不同,字节指纹也会不同,因此拒绝恢复执行。外部 API 的副作用不在这个 DB 事务之内。
下一项实验要做什么
八个步骤会连成一个可运行的成果。确认输入行的契约 → 固定源文件字节的指纹 → 保存输入和检查点 → 防止用另一个源文件恢复执行 → 以原子方式提交批次和位置 → 查询当前位置 → 按任务分开结果 → 中途失败之后安全地接着处理。
每个步骤检查的不是函数或文件是否存在,而是实际的返回值、异常和状态变化。看过正确答案之后,故意改动边界比较或清理代码,确认哪些测试会失败。说明为什么前面的测试在后面的步骤中仍然成立,并写出一个本实验不保证的运维条件。