多异步会话的原子提交实现方案?
解决异步SQLAlchemy多Session提交原子性问题
问题根源
你遇到的问题本质是:每个AsyncSession对应独立的数据库连接,各自的事务完全独立。当你分别提交两个Session时,第一个提交成功后第二个失败,就会导致数据不一致——因为两个事务没有原子性保障。
解决方案:使用PostgreSQL两阶段提交(2PC)
PostgreSQL支持两阶段提交,可以让多个独立事务实现原子性:要么全部提交成功,要么全部回滚。核心思路是先让每个事务进入"准备状态"(所有操作已执行完成,等待最终提交/回滚指令),确认所有事务都准备完成后,再统一提交;如果任何环节出错,就回滚所有准备好的事务。
代码实现
首先调整业务函数,只执行操作不提交:
async def f(session: AsyncSession): # 执行你的数据库操作,例如: session.add(YourModel(field1="value1")) # 此处不调用commit/rollback,留到后续统一处理 async def g(session: AsyncSession): # 执行另一组数据库操作,例如: session.add(AnotherModel(field2="value2"))
然后编写主逻辑,处理两阶段提交:
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine from sqlalchemy import text import asyncio import uuid async def main(): # 初始化异步引擎 engine = create_async_engine("postgresql+asyncpg://user:password@host:port/dbname") session1 = None session2 = None conn1 = None conn2 = None # 生成唯一事务ID,避免冲突 tx_id1 = f"tx_{uuid.uuid4().hex[:8]}" tx_id2 = f"tx_{uuid.uuid4().hex[:8]}" try: # 创建第一个session并执行操作,进入准备状态 conn1 = await engine.connect() session1 = AsyncSession(conn1) await f(session1) await session1.execute(text(f"PREPARE TRANSACTION '{tx_id1}'")) # 创建第二个session并执行操作,进入准备状态 conn2 = await engine.connect() session2 = AsyncSession(conn2) await g(session2) await session2.execute(text(f"PREPARE TRANSACTION '{tx_id2}'")) # 统一提交所有准备好的事务 async with engine.connect() as conn: await conn.execute(text(f"COMMIT PREPARED '{tx_id1}'")) await conn.execute(text(f"COMMIT PREPARED '{tx_id2}'")) except Exception as e: # 发生错误时,回滚所有已准备的事务 async with engine.connect() as conn: # 即使某个回滚失败,也要尝试回滚其他事务,避免悬挂事务 try: await conn.execute(text(f"ROLLBACK PREPARED '{tx_id1}'")) except: pass try: await conn.execute(text(f"ROLLBACK PREPARED '{tx_id2}'")) except: pass # 重新抛出异常,让上层感知错误 raise e finally: # 确保关闭所有session和连接 if session1: await session1.close() if conn1: await conn1.close() if session2: await session2.close() if conn2: await conn2.close() asyncio.run(main())
关键注意事项
- 事务ID唯一性:每个准备事务的ID必须唯一,这里用UUID生成短ID避免冲突,防止与其他事务ID重复导致报错。
- 权限要求:执行两阶段提交需要数据库用户拥有
PREPARE TRANSACTION权限(PostgreSQL默认允许,受限环境需提前确认)。 - 悬挂事务处理:如果系统在提交阶段崩溃,会遗留"悬挂事务",可通过
pg_prepared_xacts系统视图查看,并手动执行ROLLBACK PREPARED清理。 - 替代方案:如果业务允许,优先考虑将操作合并到同一个
AsyncSession中执行,这样天然具备原子性,无需两阶段提交——只有当必须并发执行独立操作时,才需要使用2PC。
内容的提问来源于stack exchange,提问作者John Hopfensperger
相关产品推荐
相关产品推荐

