如何配置SQLAlchemy在session.commit()执行前不与数据库通信
你观察到的现象是SQLAlchemy的默认行为,autoflush参数仅控制是否自动将ORM对象变更同步到会话操作队列,不影响session.execute()的执行时机——默认配置下每次调用session.execute()都会立即向数据库发送SQL并同步等待返回结果。
以下是两种可行的解决方案:
方案1:手动批量攒语句执行
你可以先将所有待执行的SQL、参数对暂存在内存列表中,仅在提交事务前统一调用执行,这样前面的操作仅为内存操作,耗时可以忽略。该方案不依赖特定数据库驱动,适配所有数据库类型。
from sqlalchemy import Table, Column, Integer, text from sqlalchemy import create_engine from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.orm import sessionmaker from timeit import default_timer as timer Base = declarative_base() metadata = Base.metadata tmp_table = Table( 'table_with_numbers', metadata, Column('some_number', Integer, nullable=False), ) engine=create_engine('postgresql://aaa:bbb@example.com:5432/ccc') session=sessionmaker(bind=engine, autoflush=False) with session() as s: # 初始化待执行队列 pending_ops = [] t1=timer() sql = """CREATE TEMP TABLE table_with_numbers (some_number Integer) ON COMMIT DROP;""" # 不直接执行,加入队列 pending_ops.append((text(sql), None)) t2=timer() print(f"Time taken t2-t1: {t2-t1}s") pending_ops.append((tmp_table.insert(), [ {'some_number':1}, {'some_number':2}, {'some_number':3}, ])) t3=timer() print(f"Time taken t3-t2: {t3-t2}s") # 提交前统一执行所有语句 for stmt, params in pending_ops: s.execute(stmt, params) if params else s.execute(stmt) s.commit() t4=timer() print(f"Time taken t4-t3: {t4-t3}s")
方案2:PostgreSQL流水线模式(推荐适配你的场景)
如果你使用PostgreSQL数据库,psycopg2 2.8+版本支持流水线(Pipeline)模式,开启后客户端会批量发送SQL到数据库,无需等待单条语句返回结果,最后统一处理返回值,完全符合你不需要逐次等待响应的需求,且不需要大幅修改业务代码。
from sqlalchemy import Table, Column, Integer, text from sqlalchemy import create_engine from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.orm import sessionmaker from timeit import default_timer as timer Base = declarative_base() metadata = Base.metadata tmp_table = Table( 'table_with_numbers', metadata, Column('some_number', Integer, nullable=False), ) # 正常创建引擎即可,无需额外参数 engine=create_engine('postgresql+psycopg2://aaa:bbb@example.com:5432/ccc') session=sessionmaker(bind=engine, autoflush=False) with session() as s: # 开启流水线模式 conn = s.connection() with conn.connection.pipeline(): t1=timer() sql = """CREATE TEMP TABLE table_with_numbers (some_number Integer) ON COMMIT DROP;""" s.execute(text(sql)) t2=timer() print(f"Time taken t2-t1: {t2-t1}s") s.execute(tmp_table.insert(), [ {'some_number':1}, {'some_number':2}, {'some_number':3}, ]) t3=timer() print(f"Time taken t3-t2: {t3-t2}s") s.commit() t4=timer() print(f"Time taken t4-t3: {t4-t3}s")
开启流水线后,所有execute调用只会将SQL写入客户端缓冲区,不会等待数据库返回,直到事务提交时才会同步等待所有语句执行完成,耗时表现会完全符合你的预期,且执行过程中产生的错误只会在提交阶段统一抛出。
内容的提问来源于stack exchange,提问作者Karol Zlot
相关产品推荐
相关产品推荐

