Python中Sqlite3并发如何工作?如何避免多机处理重复任务?
问题核心原因
你现在的实现有两个核心问题导致会出现重复执行:
- 你的
UPDATE语句没有加任何限定条件,执行后会把库中所有任务的状态都改成1,完全不符合单条领取的需求 - 先执行SELECT查任务、再执行UPDATE改状态的操作是非原子的,两个进程可以同时读到相同的status=0的任务,后续都执行更新就会出现重复领取
解决方案
利用SQLite的原子更新特性,把「查询待领取任务+标记为处理中」的操作合并为一步执行,从根本上避免并发争抢。
方案1(推荐,SQLite 3.35.0及以上版本支持)
SQLite 3.35.0版本开始支持RETURNING子句,可以在UPDATE执行的同时直接返回被更新的行数据,整个操作完全原子,不会出现并发冲突:
import sqlite3 def get_and_run_task(): con = sqlite3.connect("JobList.db", timeout=20) cur = con.cursor() try: # 原子操作:仅更新1条待领取任务,同时返回任务全量信息 cur.execute(""" UPDATE tasks SET status = 1 WHERE status = 0 LIMIT 1 RETURNING user, path, code, msg, report, exe, status """) task = cur.fetchone() con.commit() if not task: return user, path, code, msg, report, exe, status = task # 执行任务 run_result = dotheTask(path, exe) # 任务执行完成后更新状态和结果 cur.execute(""" UPDATE tasks SET status = 2, report = ?, msg = ? WHERE path = ? AND user = ? """, (run_result, "执行完成", path, user)) con.commit() except Exception as e: con.rollback() # 异常时可根据需求把任务状态重置为0,方便后续重试 raise e finally: con.close()
方案2(兼容旧版本SQLite)
如果你的环境SQLite版本较低不支持RETURNING,可以显式开启独占事务,同时在UPDATE时加状态校验,确保只有成功抢到更新权限的进程才会执行任务:
import sqlite3 def get_and_run_task_old(): con = sqlite3.connect("JobList.db", timeout=20) cur = con.cursor() try: # 开启独占事务,事务提交前其他进程无法修改数据库 con.execute("BEGIN EXCLUSIVE") cur.execute("SELECT user, path, code, msg, report, exe, status FROM tasks WHERE status = 0 LIMIT 1") task = cur.fetchone() if not task: con.commit() return user, path, code, msg, report, exe, status = task # 更新时再次校验状态为0,避免极端情况冲突 cur.execute("UPDATE tasks SET status = 1 WHERE status = 0 AND path = ? AND user = ?", (path, user)) con.commit() # 确认成功更新到1条记录才执行任务 if cur.rowcount <= 0: return # 执行任务 run_result = dotheTask(path, exe) # 更新最终状态 cur.execute(""" UPDATE tasks SET status = 2, report = ? WHERE path = ? AND user = ? """, (run_result, path, user)) con.commit() except Exception as e: con.rollback() raise e finally: con.close()
优化建议
- 建议给tasks表增加自增主键
id INTEGER PRIMARY KEY AUTOINCREMENT,后续定位任务用唯一id即可,不用依赖path、user等字段组合避免重复 - 可以新增
last_update_time TIMESTAMP字段,领取任务、执行中上报进度时都更新该字段,额外跑一个定时脚本把超过最大预期执行时间、状态仍为1的任务重置为0,避免机器宕机导致任务永久卡住 - 任务执行失败时可以根据业务需求增加重试次数标记,超过重试次数的任务标记为失败状态,避免无限重试
内容的提问来源于stack exchange,提问作者SharonShaji
相关产品推荐
相关产品推荐

