春江暮客

春江暮客的个人学习分享网站

用 SQLite 保存 Python 批处理进度,中断后接着运行

2026-10-09 技术
用 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,记录影响结果的输入,等结果准备好后再提交。脚本重启时就能查出已完成的工作,继续计算剩余任务。

标签

1024 12306 ablang adsense agents.md ai ai-agent ai-agents ai-seo algorithm amp antibodies automation batch-processing bioinformatics blockchain boltz bootstrapping boxes c-index cca cdn chatgpt checkpoint cli cloudflare codex cofoldarena copy cpu监控 csv cuda curl data-leakage data-processing data-quality data-validation datascience datavisualization deployment desktop-app devtools disown docker dovecot download electron esm2 esm3 esmc esmfold2 faceswap fasta fastmcp ffmpeg file-io flashppi flask folium frontend game generator git github-actions google google-research grep harness hls html http hugo indexnow javascript jev json just k-means kaggle langfuse leecode linux list litellm llm llms.txt logging logs lollipop m3u8 machine-learning macos manacher matplotlib mcp mirror model-evaluation mp3 mp4 mpnn multiomics mutation mysql nanobert nanobodies networkx nginx normalize ollama omegatherm pandas password pep-723 phaser pillow pip postfix preprocessing print protein-design protein-interactions protein-language-models protein-stability protein-structure proxy pydantic pyecharts pyqt python python3 r raincloud reproducibility requests reservoir-sampling rfoptimization rg ripgrep rosettafold3 roundcube rrsi rsi rsync s-tui sampling scale scikit-learn scrapy screen seaborn security selenium seo sha256 sklearn solana somaticsignatures spl sqlite ssh standardize static-site subprocess system-one tensorflow tkinter tron tronpy turtle typesafe-ai usdt uv vhhbert vite webp wordcloud wordpress workflow yaml 后台 寓言 概率 经济 贸易 迅雷解析 钱包

友情链接

其它