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会为每个线程创建独立会话保证线程安全,但还是出现了锁错误,请问问题出在哪?怎么解决?
原因分析
- SQLite锁机制限制:SQLite是文件型数据库,写操作会独占整个数据库锁,多线程同时执行写操作时,后续线程会因获取锁失败抛出锁错误。即使每个线程有独立会话,底层数据库文件的锁冲突依然存在。
- 重复执行DDL操作:每个任务都调用
Base.metadata.create_all(...),多线程同时执行DDL会加剧锁竞争——DDL操作需要更高级别的数据库锁。 - 线程池并发过高:线程数与任务数一致导致所有任务同时启动,并发执行数据库写操作,直接触发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
相关产品推荐
相关产品推荐

