Pytest+SqlAlchemy+FastAPI测试问题:同步插入数据异步无法读取
问题:同步会话插入数据后异步API无法查询到数据
问题背景
- 应用分为两部分:提供数据查询的API(使用异步数据库连接)、消费PubSub事件写入数据的同步消费者(使用同步数据库连接)
- 生产环境代码运行正常,但测试场景下:用同步会话插入数据后,通过异步客户端调用API无法获取到数据
- 排查细节:同步会话插入的数据已写入数据库(通过pdb验证),但异步会话/客户端查询结果为空,怀疑测试fixture配置存在问题
补充疑问
是否可以将同步引擎插入数据后创建的保存点传递给异步fixture,让异步引擎从该保存点开始运行?
尝试过的方案
参考了SQLAlchemy「将会话加入外部事务」的方案,配置了同步/异步的测试会话和客户端,但问题仍未解决。
相关代码
conftest.py
import sqlalchemy as sa from fastapi.testclient import TestClient from sqlalchemy.orm import sessionmaker from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from httpx import AsyncClient import pytest @pytest.fixture(scope="session") def engine(config: TestConfig): return sa.create_engine(config.DATABASE_URL) @pytest.fixture(scope="session") def async_engine(config: TestConfig): return create_async_engine( config.ASYNC_DATABASE_URL, pool_size=10, echo=True, max_overflow=10 ) @pytest.fixture(scope="session") def TestingSessionLocal(engine): return sessionmaker( autocommit=False, autoflush=False, bind=engine, ) @pytest.fixture(scope="session") def TestingSessionLocalAsync(async_engine): return sessionmaker( bind=async_engine, autoflush=False, future=True, class_=AsyncSession, expire_on_commit=False, ) @pytest.fixture def db(engine, TestingSessionLocal): connection = engine.connect() transaction = connection.begin() session = TestingSessionLocal(bind=connection) session.begin_nested() @sa.event.listens_for(session, "after_transaction_end") def end_savepoint(session, transaction): if transaction.nested and not transaction._parent.nested: session.begin_nested() yield session session.close() transaction.rollback() connection.close() @pytest.fixture async def async_db(async_engine, TestingSessionLocalAsync): async with async_engine.connect() as conn: await conn.begin() await conn.begin_nested() async_session = TestingSessionLocalAsync() @sa.event.listens_for(async_session.sync_session, "after_transaction_end") def end_savepoint(session, transaction): if conn.closed: return if not conn.in_nested_transaction(): conn.sync_connection.begin_nested() yield async_session await async_session.close() await conn.close() @pytest.fixture def app(config): return create_app(config) @pytest.fixture def client(app, db): def override_get_db(): yield db app.dependency_overrides[get_db] = override_get_db yield TestClient(app) del app.dependency_overrides[get_db] @pytest.fixture() async def async_client(app, async_db): def override_get_async_db(): yield async_db app.dependency_overrides[get_async_db] = override_get_async_db async with AsyncClient( app=app, base_url="http://localhost:8080", headers={"Content-Type": "application/json"}, ) as client: yield client del app.dependency_overrides[get_async_db]
test.py
@pytest.mark.asyncio async def test_thing(async_client, db: Session): thing = Thing() db.add(thing) db.flush() db.refresh(thing) # db.commit() # 尝试添加此代码后结果无变化 res1 = await async_db.execute(select(thing.name).select_from(thing)) res2 = db.execute(select(thing.name).select_from(thing)) pdb.set_trace() get_thing = await async_client.get( f"/thing/{thing.id}", ) assert len(get_thing) == 1
pdb调试结果

解决方案
核心问题分析
当前测试中同步会话和异步会话使用独立的数据库连接与事务上下文,同步会话的嵌套事务(保存点)无法被异步会话共享,导致异步会话看不到同步插入的数据。
具体修改步骤
1. 修改async_db fixture,共享同步会话的事务上下文
让异步会话绑定到同步会话的连接上,确保两者处于同一个事务中:
# 更新conftest.py中的async_db fixture @pytest.fixture async def async_db(db, TestingSessionLocalAsync): # 复用同步会话的底层连接创建异步连接 sync_conn = db.connection() async_conn = await TestingSessionLocalAsync.bind.connect( connection=sync_conn.connection, reuse_connection=True ) await async_conn.begin_nested() async_session = TestingSessionLocalAsync(bind=async_conn) @sa.event.listens_for(async_session.sync_session, "after_transaction_end") def end_savepoint(session, transaction): if async_conn.closed: return if not async_conn.in_nested_transaction(): async_conn.sync_connection.begin_nested() yield async_session await async_session.close() await async_conn.close()
2. 调整测试用例,确保变更写入事务
同步会话插入数据后,提交到嵌套事务(保存点),确保变更对共享上下文可见:
# 更新test.py的测试用例 @pytest.mark.asyncio async def test_thing(async_client, db: Session): thing = Thing() db.add(thing) db.commit() # 提交到嵌套事务,确保变更持久化到当前事务上下文 db.refresh(thing) get_thing = await async_client.get(f"/thing/{thing.id}") # 注意:get_thing是响应对象,需取JSON内容判断长度 assert len(get_thing.json()) == 1
关键说明
- 同步与异步会话必须共享同一个事务上下文,才能互相可见未提交的变更(符合数据库ACID特性)
- 使用
reuse_connection=True让异步连接复用同步连接的底层数据库连接,保证事务上下文一致 - 测试结束后,db fixture的回滚操作会同时撤销同步和异步会话的所有变更,保证测试隔离性
内容的提问来源于stack exchange,提问作者fallen
相关产品推荐
相关产品推荐

