在后台服务与数据管线中,Python常被用来编写定时或手动触发的批处理任务,例如对账、日志清洗、报表生成等。这类任务往往需要处理成千上万条记录,执行时间从几分钟到数小时不等。如果在运行过程中遇到进程被杀死、机器重启、网络抖动或者主动暂停,再次启动同一个脚本时,如果没有特殊设计,就会从头开始处理,已经写入数据库的结果可能被重复写入,造成数据不一致甚至资损。

可重入设计指的是一个计算过程在任意时刻被中断,之后能够重新进入并从合理的进度继续,且最终结果与一次性跑完完全一致。它不等同于简单的重试,重试可能放大错误,而可重入强调过程的可恢复与结果的幂等。下面从几个关键技术角度展开说明。
状态持久化与断点记录
实现可重入的第一一步是把任务的执行进度落到可靠的外部存储中,而不是只存在于内存变量里。最常见的做法是利用数据库表或者独立的进度文件记录“已处理到哪个游标”。例如对订单表按主键分批,每次提交一批后把最大主键写入进度表。进程重启后先读进度,再从该主键之后查询,自然跳过了历史数据。
如果任务本身没有数据库依赖,也可以使用本地文件或对象存储保存JSON格式的断点信息。需要注意的是,写入进度必须与业务数据处理放在同一个事务或至少保证顺序一致,否则可能出现“数据已处理但进度未记”或反过来“进度记了但数据回滚”的缝隙。下面给出一个基于SQLite的简易进度管理示例:
import sqlite3
def init_progress(conn):
conn.execute('''
CREATE TABLE IF NOT EXISTS batch_progress (
task_name TEXT PRIMARY KEY,
last_id INTEGER NOT NULL
)
''')
conn.commit()
def get_last_id(conn, task_name):
row = conn.execute(
'SELECT last_id FROM batch_progress WHERE task_name=?',
(task_name,)
).fetchone()
return row[0] if row else 0
def update_last_id(conn, task_name, last_id):
conn.execute('''
INSERT INTO batch_progress(task_name, last_id) VALUES(?, ?)
ON CONFLICT(task_name) DO UPDATE SET last_id=excluded.last_id
''', (task_name, last_id))
conn.commit()
def process_batch():
conn = sqlite3.connect('batch.db')
init_progress(conn)
task = 'order_sync'
start = get_last_id(conn, task)
rows = conn.execute(
'SELECT id, amount FROM orders WHERE id > ? ORDER BY id LIMIT 500',
(start,)
).fetchall()
for r in rows:
# 模拟业务处理
print('handle order', r[0])
start = r[0]
update_last_id(conn, task, start)
conn.close()
process_batch()
上面的代码每次最多取五百条,处理完整体批次后才更新last_id。如果中途崩溃,下次启动会从已保存的last_id继续,从而避免重复。实际生产可以把批次缩小并使用每批独立事务,进一步降低丢失窗口。对于分布式多实例场景,进度表应加上实例标识或采用分布式锁,防止相互覆盖。
幂等性与数据冲突控制
仅有断点还不够,因为断点更新本身可能失败,或者一个批次内处理了部分记录后崩溃,导致这批记录处于“半完成”状态。此时重入会再次触碰这些记录,如果写操作不是幂等的,就会出错。幂等意味着同一输入多次执行效果相同,比如“设置状态为已同步”是幂等的,而“余额加十元”不是。
在Python批处理中,常用数据库唯一约束、乐观锁或条件更新来保证幂等。举例来说,把处理状态字段加上CHECK约束,更新时写UPDATE orders SET synced=1 WHERE id=? AND synced=0,受影响行数为零就说明之前已处理。另一种方式是使用版本号,每次更新附带原版本,冲突则跳过。如下示例展示带版本号的更新:
import sqlite3
def safe_update(conn, order_id, old_version):
cur = conn.execute('''
UPDATE orders SET synced=1, version=version+1
WHERE id=? AND version=?
''', (order_id, old_version))
conn.commit()
return cur.rowcount
conn = sqlite3.connect('batch.db')
# 假设读取时拿到版本
print(safe_update(conn, 1001, 1))
conn.close()
当返回值为零,说明版本不匹配,大概率已被其他重入流程处理过,此时直接忽略即可。通过这样的冲突控制,即便断点粒度较粗,也不会产生重复副作用。还需要注意外部接口调用,如发送通知或请求第三方API,应记录请求指纹并先查重,否则网络超时引起的重发也会制造垃圾数据。
任务调度与异常隔离策略
可重入设计最后要落地到运行方式上。很多团队用cron或APScheduler拉起Python脚本,如果上一次还没跑完就又触发,会造成并发写进度。解决办法是在任务入口用文件锁或数据库锁互斥,或者让调度器只负责投递信号,真正执行由常驻worker取走。以下示例用fcntl做简单单实例约束:
import fcntl
import sys
lock_file = open('/tmp/batch.lock', 'w')
try:
fcntl.flock(lock_file, fcntl.LOCK_EX | fcntl.LOCK_NB)
except BlockingIOError:
print('another instance is running')
sys.exit(0)
# 此处调用可重入批处理主逻辑
print('start reentrant batch')
除了互斥,还要对单条记录错误进行隔离。不要因为一行脏数据就抛异常退出导致整个任务不可重入,应当捕获异常、写入死信表并继续下一条。这样重入时只补跑死信,而不必重扫全量。配合前面说的进度与幂等,整个系统就能在频繁中断的环境中保持数据准确。
综合来看,Python批处理的可重入并不是某个单一技巧,而是状态外置、写操作幂等以及运行隔离三者的组合。前期多花少量代码成本,后期能省掉大量人工修数和故障排查时间,对中大型数据任务尤其值得投入。