使用SQLAlchemy的DB Session并行调用报错问题排查
问题分析与解决方案
问题场景
当前创建了全局数据库会话:
db_engine = create_engine(DB_URL, pool_pre_ping=True, echo=False, pool_size=POOL_SIZE, max_overflow=MAX_OVERFLOW) Base.metadata.create_all(db_engine) db_session = sessionmaker(bind=db_engine)()
并在Flask路由中直接引用该全局会话:
from db_settings import db_session @application.route('/') def my_home(): try: my_obj = MyObj("dummy_fname", "dummy_lname") db_session.merge(my_obj) db_session.commit() except exc.IntegrityError: db_session.rollback() finally: db_session.close()
测试表现:串行调用路由可成功创建5条记录,但多线程并行调用时仅成功1条,其余报错。测试代码如下:
import requests import threading no_parallel_hits = no_of_sequential_hits = 5 def send_request(): url = 'http://127.0.0.1:5000/' response = requests.get(url) print(response.text) def main_parallel(): threads = [] for _ in range(no_parallel_hits): thread = threading.Thread(target=send_request) threads.append(thread) thread.start() for thread in threads: thread.join() def main_sequential(): for _ in range(no_of_sequential_hits): send_request() if __name__ == "__main__": main_parallel() # main_sequential()
核心原因
问题根源并非db_session.close()的影响,而是全局db_session是单例实例,SQLAlchemy的Session本身不具备线程安全性:
- 串行调用时,请求依次执行,前一个请求关闭session后,下一个请求可重新激活该实例,无并发冲突;
- 并行调用时,多个线程同时操作同一个session实例:
- 一个线程执行commit/rollback时,其他线程可能正在修改session状态;
- 某线程执行close后,其他线程的session操作会因连接关闭报错;
- 并发的merge、commit操作会导致事务状态混乱,最终仅一个线程的操作能成功,其余因session状态异常或事务冲突失败。
解决方案
方案1:每个请求创建独立Session
在路由内部创建新的session实例,确保每个请求拥有独立会话:
from db_settings import sessionmaker, db_engine @application.route('/') def my_home(): db_session = sessionmaker(bind=db_engine)() try: my_obj = MyObj("dummy_fname", "dummy_lname") db_session.merge(my_obj) db_session.commit() except exc.IntegrityError: db_session.rollback() finally: db_session.close()
方案2:使用Scoped Session(推荐)
SQLAlchemy的scoped_session会为每个线程维护独立的Session实例,适配Web场景:
修改db_settings.py:
from sqlalchemy.orm import scoped_session, sessionmaker db_engine = create_engine(DB_URL, pool_pre_ping=True, echo=False, pool_size=POOL_SIZE, max_overflow=MAX_OVERFLOW) Base.metadata.create_all(db_engine) db_session = scoped_session(sessionmaker(bind=db_engine))
路由中使用并清理当前线程的session实例:
from db_settings import db_session from sqlalchemy import exc @application.route('/') def my_home(): try: my_obj = MyObj("dummy_fname", "dummy_lname") db_session.merge(my_obj) db_session.commit() except exc.IntegrityError: db_session.rollback() finally: db_session.remove() # 清理当前线程的session实例,而非直接close
scoped_session会自动为每个线程分配独立Session,remove()操作仅清理当前线程的实例,避免线程间状态污染。
内容的提问来源于stack exchange,提问作者abhi
相关产品推荐
相关产品推荐

