用 SQLite 保存 Python 批处理进度,中断后接着运行
批处理脚本已经完成几百条记录,中途退出后,下次运行却从头开始。用一个小型 SQLite 数据库保存已完成结果,就能让脚本重启后接着处理剩余任务。
本文实现一个单进程批处理工具。每完成一个任务,就把结果提交为一行记录。重启时先检查当前输入,再跳过已保存且匹配的任务。示例只计算文本长度,不需要模型、API 密钥或外部程序,可以直接验证恢复过程。
1. 把结果和完成状态放在一起
如果进度文件已经写入 done=true,结果文件却还没保存,重启后就很难判断该不该重做。这里直接把结果行作为完成记录:
check input fingerprint → calculate → serialize result → insert row → commit → report DONE
表里只有三列:稳定的任务 ID、输入与计算版本的指纹,以及 JSON 结果。任务 ID 是主键。同一个 ID 换了文本,或者计算版本发生变化,都不能继续使用原来的结果。
这个示例假设只有一个进程使用数据库,数据库放在适合 SQLite 的本地文件系统上,结果也保存在数据库内部。大型输出文件、多进程计算和 API 操作需要额外协调。
2. 下载脚本,主动暂停一次
使用 Python 3.12 或更新版本,解释器需要提供标准库 sqlite3 模块。示例已在 macOS、Python 3.14.7 上测试。下面的命令适用于 bash 或 zsh。
把 resume_batch.py 下载到新的工作目录,再把以下内容保存为 jobs.json:
[
{"id": "seq_001", "text": "MKT"},
{"id": "seq_002", "text": "ACDE"},
{"id": "seq_003", "text": "GGGGG"}
]
运行:
python3 --version
python3 resume_batch.py --help
mkdir sqlite-resume-demo
python3 resume_batch.py jobs.json sqlite-resume-demo/progress.sqlite --stop-after 2
批处理命令的预期输出:
DONE seq_001
DONE seq_002
STOPPED: new=2; skipped=0
--stop-after 表示新增并提交指定数量的结果后暂停,退出码为 0。跳过的任务不计入这个数量。这是可重复的主动暂停演示,不是断电模拟。重新做完整演示时,请使用新的目录。
3. 继续运行,再重复一次
使用相同的输入文件和数据库,去掉暂停参数:
python3 resume_batch.py jobs.json sqlite-resume-demo/progress.sqlite
前两个任务已经有结果:
SKIP seq_001
SKIP seq_002
DONE seq_003
FINISHED: new=1; skipped=2
再运行一次同样的命令,会输出三行 SKIP,最后是 FINISHED: new=0; skipped=3。
脚本在数据库提交成功后才打印 DONE。进程可能在提交和打印之间退出,所以终端少了一行消息,并不代表结果没有保存。恢复判断以数据库为准。
4. 完整实现
"""Resume a single-worker text batch using SQLite. Python 3.12+."""
import argparse
import hashlib
import json
from pathlib import Path
import sqlite3
import sys
WORKER_VERSION = "text-length-v1"
def load_jobs(path):
jobs = json.loads(Path(path).read_text(encoding="utf-8"))
if not isinstance(jobs, list):
raise ValueError("input must be a JSON list")
seen = set()
for job in jobs:
if not isinstance(job, dict) or set(job) != {"id", "text"}:
raise ValueError("each job needs exactly id and text")
job_id, text = job["id"], job["text"]
if not isinstance(job_id, str) or not job_id.strip():
raise ValueError("job id must be a nonempty string")
if not isinstance(text, str):
raise ValueError("job text must be a string")
if job_id in seen:
raise ValueError(f"duplicate job id: {job_id}")
seen.add(job_id)
return jobs
def fingerprint(text):
payload = json.dumps([WORKER_VERSION, text], ensure_ascii=True)
return hashlib.sha256(payload.encode("utf-8")).hexdigest()
def calculate(text):
return {"length": len(text)}
def run_batch(jobs, database, *, stop_after=None):
completed = skipped = 0
con = sqlite3.connect(database, autocommit=False)
try:
with con:
con.execute("""CREATE TABLE IF NOT EXISTS results (
job_id TEXT PRIMARY KEY,
fingerprint TEXT NOT NULL,
result_json TEXT NOT NULL
)""")
# Check every current input before calculating any new result.
for job in jobs:
old = con.execute(
"SELECT fingerprint FROM results WHERE job_id = ?",
(job["id"],),
).fetchone()
if old is not None and old[0] != fingerprint(job["text"]):
raise ValueError(f"input or worker changed for {job['id']}; use a new database")
for job in jobs:
old = con.execute(
"SELECT 1 FROM results WHERE job_id = ?", (job["id"],)
).fetchone()
if old is not None:
skipped += 1
print(f"SKIP {job['id']}", flush=True)
continue
result = calculate(job["text"])
encoded = json.dumps(result, allow_nan=False, sort_keys=True)
# The result row itself is the completion record.
with con:
con.execute(
"INSERT INTO results VALUES (?, ?, ?)",
(job["id"], fingerprint(job["text"]), encoded),
)
completed += 1
print(f"DONE {job['id']}", flush=True)
if stop_after is not None and completed >= stop_after:
print(f"STOPPED: new={completed}; skipped={skipped}", flush=True)
return
print(f"FINISHED: new={completed}; skipped={skipped}", flush=True)
finally:
con.close()
def positive_int(value):
number = int(value)
if number < 1:
raise argparse.ArgumentTypeError("must be at least 1")
return number
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("input", type=Path, help="JSON list of id/text jobs")
parser.add_argument("database", type=Path, help="SQLite checkpoint; parent must exist")
parser.add_argument("--stop-after", type=positive_int, help="pause after N new results")
args = parser.parse_args()
try:
jobs = load_jobs(args.input)
run_batch(jobs, args.database, stop_after=args.stop_after)
except (OSError, ValueError, sqlite3.Error) as exc:
print(f"ERROR: {exc}", file=sys.stderr)
return 1
except KeyboardInterrupt:
print("INTERRUPTED: rerun with the same input and database", file=sys.stderr)
return 130
return 0
if __name__ == "__main__":
raise SystemExit(main())
命令行入口会在打开数据库之前校验整个 JSON 列表。ID 必须是非空字符串,且不能重复;文本必须是字符串。空文本合法,长度为零。SQL 通过占位符传入数据,不把输入拼进 SQL 语句。
指纹包含 WORKER_VERSION 和完整文本。输入顺序可以变化,但已有 ID 的文本改变后,需要换一个数据库。指纹不会自动发现 Python 代码改动:计算逻辑变了,就要更新版本。实际推理流程还应把模型标识、参数和影响结果的依赖纳入任务身份。
代码明确设置 autocommit=False。Python 的事务控制文档说明了这个模式。连接上下文管理器会提交成功的代码块,并回滚失败的代码块,但不会关闭连接,因此还需要 finally: con.close()。
5. 查看数据库中的结果
把以下代码保存为 inspect_results.py,放在下载脚本旁边:
import json
import sqlite3
con = sqlite3.connect("sqlite-resume-demo/progress.sqlite")
try:
for job_id, result_json in con.execute(
"SELECT job_id, result_json FROM results ORDER BY job_id"
):
print(job_id, json.loads(result_json)["length"])
finally:
con.close()
运行 python3 inspect_results.py:
seq_001 3
seq_002 4
seq_003 5
示例统计 Python 字符串的字符数量,不是 UTF-8 字节数,也不是屏幕上可见的字形数量。它不校验蛋白质序列。接入自己的流程时,替换 calculate(),检查计算结果,再更新 WORKER_VERSION。
从输入列表移除一个任务,不会删除数据库中的旧结果。上面的查询会显示所有已保存记录;导出某个输入文件的结果时,应按当前任务 ID 筛选。
6. 验证输入改变后会被拒绝
把以下内容保存为 changed.json:
[
{"id": "seq_001", "text": "MKA"},
{"id": "seq_002", "text": "ACDE"},
{"id": "seq_003", "text": "GGGGG"}
]
运行:
python3 resume_batch.py changed.json sqlite-resume-demo/progress.sqlite
脚本返回退出码 1,并报告:
ERROR: input or worker changed for seq_001; use a new database
两段文本长度都是三,但内容不同。脚本会在计算任何新增任务前发现这个变化,原来的结果仍然保留。要处理新数据,可以选择 sqlite-resume-demo/changed.sqlite 这样的新路径。
数据集校验清单可以标识输入背后的整套文件。这里的单任务指纹则判断某条已保存结果是否属于当前任务定义。
7. 明确恢复边界
| 退出位置 | 下次运行的行为 |
|---|---|
| 计算尚未完成 | 重新计算该任务。 |
| 计算完成,但还没有成功提交 | 如果没有已提交的结果行,就重新计算。 |
已提交,还没打印 DONE |
找到结果行,跳过计算。 |
| 后续任务失败 | 保留之前已提交的结果;重启后重做未完成任务。 |
| 已有任务的输入或计算版本变化 | 在预检查阶段退出,需要使用新数据库。 |
本文验证了暂停后继续运行、重复跳过、同长度文本变化、计算版本变化、不合法输入、Unicode 文本、数据库插入失败,以及第三条结果保存前的进程突然退出。这些检查验证了进程恢复,没有进行断电测试。
SQLite 的原子提交说明介绍了事务恢复,以及它依赖的文件系统与存储条件。数据库应放在合适的存储上,保留其日志文件;调整持久性设置前,需要了解影响。
如果计算函数调用 API 或写入另一个文件,进程可能在外部操作成功后、数据库记录之前退出。下次运行会再次执行该操作。可以使用服务支持的幂等键,或者让重复执行同一任务仍然安全。这个示例不保证外部操作只发生一次。
执行外部程序时,可以配合子进程超时与日志工具,先验证输出,再提交结果。对于独立的小型状态文件,原子写入 JSON也是可搭配使用的方法。
8. 常见错误与处理
| 错误 | 处理方法 |
|---|---|
duplicate job id |
为每条输入分配唯一且稳定的 ID。 |
input or worker changed |
任务定义改变后,选择新的数据库路径。 |
unable to open database file |
创建父目录,并检查写入权限。 |
database is locked |
停止重叠运行的进程;本示例只支持一个计算进程。 |
| 修改代码后结果不符合预期 | 更新 WORKER_VERSION,并使用新数据库。 |
9. 小结
先确定稳定的任务 ID,记录影响结果的输入,等结果准备好后再提交。脚本重启时就能查出已完成的工作,继续计算剩余任务。
- 原文作者:春江暮客
- 原文链接:https://www.bobobk.com/python-sqlite-resume-batch.html
- 版权声明:本作品采用 知识共享署名-非商业性使用-禁止演绎 4.0 国际许可协议 进行许可,非商业转载请注明出处(作者,原文链接),商业转载请联系作者获得授权。