如何使用SQLAlchemy Asyncio插入记录并获取主键?验证并发安全性
问题描述
我在Python(FastAPI)项目中需要向PostgreSQL(或任意数据库)插入记录并获取插入的主键,查阅SQLAlchemy AsyncIO官方文档后没找到易懂的基于会话的解决方案。目前少量用户使用代码正常,但想确认这套代码在多并发用户/请求场景下有没有bug,以及高负载环境下的表现。
当前代码实现
数据库引擎创建
async def create_postgresql_engine(): engine = create_async_engine(ConnectionString, future=True, echo=True) return engine
业务接口入口
async def post_apps_logs_def( arg1 ,arg2): return await post_apps_logs(engine=await create_postgresql_engine(),arg1, arg2)
日志插入逻辑
async def post_apps_logs(engine, arg1, arg2): result = pd.DataFrame() try: stmt = APP_LOGS( arg1=arg1, arg2=arg2 ) async_session = async_sessionmaker(engine, expire_on_commit=False) result=await insert_stmt(stmt,async_session) except SQLAlchemyError as e: error = str(e.__cause__) raise RuntimeError(error) from e return result
通用插入语句
async def insert_stmt(stmt,async_session: async_sessionmaker[AsyncSession]): df=pd.DataFrame({}) async with async_session() as session: try: async with session.begin(): session.add(stmt) await session.flush() await session.refresh(stmt) df= (stmt.PK) except SQLAlchemyError as e: error = str(e.__cause__) await session.rollback() raise RuntimeError(error) from e finally: await session.close() return df
代码问题分析
- 重复创建引擎与会话工厂:每次请求都会创建新的数据库引擎和会话工厂,SQLAlchemy的Engine是协程安全的全局资源,重复创建会浪费连接池资源,高负载下会导致性能急剧下降。
- 冗余的会话与事务处理:
async with async_session()上下文管理器会自动关闭会话,finally中的await session.close()属于多余操作,甚至可能引发上下文冲突;async with session.begin()已经是事务上下文,异常时会自动回滚,手动调用await session.rollback()会导致重复回滚。 - 变量类型混淆:
df初始化为DataFrame,但最终返回单个主键值,变量命名容易造成误解,且完全没必要引入pandas处理单个值。 - 错误处理冗余:
post_apps_logs和insert_stmt重复捕获SQLAlchemyError并封装,导致错误栈冗余。
优化后的代码实现
全局资源初始化(建议在项目启动时执行)
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import async_sessionmaker from sqlalchemy.exc import SQLAlchemyError # 全局单例数据库引擎,仅初始化一次 engine = create_async_engine(ConnectionString, future=True, echo=True) # 全局会话工厂,复用该实例创建会话 AsyncSessionLocal = async_sessionmaker(engine, expire_on_commit=False)
通用插入函数
async def insert_model(model_instance) -> int: async with AsyncSessionLocal() as session: async with session.begin(): session.add(model_instance) await session.flush() # 触发数据库插入,生成主键 await session.refresh(model_instance) # 同步实例的主键属性 return model_instance.PK
业务插入逻辑
async def post_apps_logs(arg1, arg2) -> int: log_instance = APP_LOGS(arg1=arg1, arg2=arg2) try: return await insert_model(log_instance) except SQLAlchemyError as e: raise RuntimeError(f"插入日志失败: {str(e.__cause__)}") from e async def post_apps_logs_def(arg1, arg2) -> int: return await post_apps_logs(arg1, arg2)
并发与高负载场景验证
- 会话安全性:SQLAlchemy的
AsyncSession是协程隔离的,每个请求使用独立的会话实例,多并发下不会出现数据混淆或竞争问题。 - 主键获取可靠性:
flush()会将插入操作发送到数据库并生成主键,refresh()会将数据库生成的主键同步到模型实例,该流程在PostgreSQL(自增主键/序列主键)及其他支持返回插入主键的数据库中完全可靠。 - 连接池性能:使用全局单例Engine后,SQLAlchemy会自动维护连接池(默认大小为5),高负载下可通过调整
pool_size和max_overflow参数优化连接池配置,避免连接耗尽。
内容的提问来源于stack exchange,提问作者Moh-Spark
相关产品推荐
相关产品推荐

