You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多异步会话的原子提交实现方案?

解决异步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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 22:52:43