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

APScheduler多线程结合SQLAlchemy遇SQLite锁错误求解决

问题:多线程SQLite操作出现database is locked错误

我在APScheduler的BlockingScheduler中添加了多个任务,使用ThreadPoolExecutor且线程数等于任务数,所有任务通过SQLAlchemy操作同一个SQLite数据库时,出现以下错误:

sqlalchemy.exc.OperationalError: (sqlite3.OperationalError) database is locked

我的SQLAlchemy基础配置已使用scoped_session和sessionmaker:

from os.path import join, realpath
from sqlalchemy import create_engine
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import scoped_session, sessionmaker
from os import environ

db_name = environ.get("DB_NAME")
db_path = realpath(join("data", db_name))
engine = create_engine(f"sqlite:///{db_path}", pool_pre_ping=True)
session_factory = sessionmaker(bind=engine)
Session = scoped_session(session_factory)
Base = declarative_base()

调度任务类示例:

from os.path import realpath
from os import environ
from typing import List
from app.data_structures.base import Base, Session, engine
from app.data_structures.job import Job
from app.data_structures.scheduled_job import ScheduledJob
from app.data_structures.user import User
from app.resources import AppResources


class AccountingProcessorJob(ScheduledJob):
    name: str = "Accounting Processor"
    def __init__(self, resources: AppResources, depends: List[str] = None) -> None:
        super().__init__(resources)

    def job_function(self) -> None:
        account_dir = realpath(environ.get("ACCOUNTING_DIRECTORY"))
        Base.metadata.create_all(engine, Base.metadata.tables.values(), checkfirst=True)
        session = Session()
        try:
            # 示例操作
            user_name = "test_user"
            jobs = []
            user = User(user_name=user_name)
            session.add(user)
            user.jobs.extend(jobs)
            session.commit()
        except:
            session.rollback()
        finally:
            Session.remove() 

ORM模型示例(User类):

from sqlalchemy import Column, Integer, String
from sqlalchemy.orm import Mapped, mapped_column, relationship
from app.data_structures.base import Base
from app.data_structures.job import Job


class User(Base):
    __tablename__ = "users"

    user_name: Mapped[str] = mapped_column(primary_key=True)
    employee_number = Column(Integer)
    manager = relationship("User", remote_side=[user_name], post_update=True)
    jobs: Mapped[list[Job]] = relationship()

    def __init__(
        self,
        user_name: str,
        employee_number: int = None,
        manager: str = None,
    ) -> None:
        self.user_name = user_name
        self.employee_number = employee_number
        self.manager = manager

我原本以为scoped_session会为每个线程创建独立会话保证线程安全,但还是出现了锁错误,请问问题出在哪?怎么解决?


原因分析
  1. SQLite锁机制限制:SQLite是文件型数据库,写操作会独占整个数据库锁,多线程同时执行写操作时,后续线程会因获取锁失败抛出锁错误。即使每个线程有独立会话,底层数据库文件的锁冲突依然存在。
  2. 重复执行DDL操作:每个任务都调用Base.metadata.create_all(...),多线程同时执行DDL会加剧锁竞争——DDL操作需要更高级别的数据库锁。
  3. 线程池并发过高:线程数与任务数一致导致所有任务同时启动,并发执行数据库写操作,直接触发SQLite的锁冲突。

解决方案

1. 优化SQLite连接配置

创建engine时添加SQLite专属参数,提升锁等待超时时间并适配多线程环境:

engine = create_engine(
    f"sqlite:///{db_path}",
    pool_pre_ping=True,
    connect_args={
        "check_same_thread": False,  # 配合scoped_session使用,允许连接跨线程共享
        "timeout": 30  # 锁等待超时时间从默认5秒调至30秒,降低锁冲突报错概率
    }
)

2. 提前执行DDL操作

把Base.metadata.create_all(...)从任务函数中移除,放到程序启动的初始化阶段,仅执行一次:

# 程序启动时执行(比如主函数开头)
Base.metadata.create_all(engine, checkfirst=True)

3. 控制线程池并发数

SQLite写操作无法真正并发,建议将ThreadPoolExecutor的线程数设为1(若任务以写为主):

from apscheduler.executors.pool import ThreadPoolExecutor

executors = {
    'default': ThreadPoolExecutor(1)  # 单线程执行写任务,避免锁冲突
}
scheduler = BlockingScheduler(executors=executors)

若任务包含读写分离场景,可单独为读任务配置多线程,写任务用单线程队列执行。

4. 保持会话正确清理

你的代码中已在finally块调用Session.remove(),这部分是正确的——scoped_session会将会话绑定到线程本地存储,确保每个线程获取独立会话,继续保持该逻辑即可。


内容的提问来源于stack exchange,提问作者abinitio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 19:05:32