You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python中Sqlite3并发如何工作?如何避免多机处理重复任务?

问题核心原因

你现在的实现有两个核心问题导致会出现重复执行:

  1. 你的UPDATE语句没有加任何限定条件,执行后会把库中所有任务的状态都改成1,完全不符合单条领取的需求
  2. 先执行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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.06 01:24:00