如何用asyncpg实现PostgreSQL异步Mock并解决FastAPI多协程测试问题
解决asyncpg多协程测试中的连接冲突问题
问题背景
我们维护一个包含HTTP和WebSocket端点的FastAPI服务,使用PostgreSQL+asyncpg实现CRUD操作。其中POST /todo端点在插入数据后,会启动异步任务通知WebSocket监听者,核心代码如下:
async def notify_todo_listeners(): todos = await db_handler.fetch_todos() notify_all(todos) @app.post("/todo") async def post_todo(request: Request): todo = await db_handler.insert(await request.json()) asyncio.create_task(notify_todo_listeners()) return todo
为了测试该端点,我们用Docker创建临时PostgreSQL数据库,通过pytest fixture创建数据库连接并开启事务,测试后回滚所有变更,核心测试代码如下:
@pytest_asyncio.fixture(scope="function") async def session(monkeypatch): connection = await asyncpg.connect(CONNECTION_STRING) transaction = connection.transaction() await transaction.start() async def mock_get_connection(): return connection monkeypatch.setattr(database_handler, "get_connection", mock_get_connection) yield connection await transaction.rollback() await connection.close() async def test_post_todo(session): async with AsyncClient(app=app, base_url="http://test") as client: response = await client.post("/todo", json={"content": "test task"}) assert response.status_code == 200 # 其他断言逻辑
测试时触发错误:
exception=InterfaceError('cannot perform operation: another operation is in progress')
错误原因
asyncpg的单个连接是单协程独占资源,同一时间只能处理一个异步操作。测试中,POST /todo协程正在使用连接执行插入操作,同时notify_todo_listeners异步任务又尝试复用同一个连接执行查询,导致连接冲突报错。
解决方法
方案1:使用连接池+共享事务上下文
通过连接池为每个协程提供独立连接,同时所有连接绑定到同一个顶级事务,确保测试结束后可一次性回滚所有变更。
修改后的测试fixture:
import asyncpg import pytest_asyncio from httpx import AsyncClient @pytest_asyncio.fixture(scope="function") async def db_pool(monkeypatch): # 创建连接池 pool = await asyncpg.create_pool(CONNECTION_STRING) # 获取顶级连接并启动事务 async with pool.acquire() as top_conn: top_transaction = top_conn.transaction() await top_transaction.start() async def mock_get_connection(): # 为每个协程从池获取独立连接,并启动嵌套事务 conn = await pool.acquire() await conn.transaction(start=True) return conn monkeypatch.setattr(database_handler, "get_connection", mock_get_connection) yield pool # 回滚顶级事务,所有嵌套事务自动回滚 await top_transaction.rollback() # 测试结束后关闭连接池 await pool.close() async def test_post_todo(db_pool): async with AsyncClient(app=app, base_url="http://test") as client: response = await client.post("/todo", json={"content": "test todo"}) assert response.status_code == 200 # 补充其他断言,比如验证WebSocket通知逻辑等
关键注意事项
- 数据库操作必须正确释放连接:确保
db_handler中的方法使用async with管理连接,比如:async def fetch_todos(): async with await get_connection() as conn: return await conn.fetch("SELECT * FROM todos") - 避免遗漏await:原代码中如果
db_handler.insert或fetch_todos是异步方法,必须添加await,否则会返回协程对象导致逻辑错误。 - 事务嵌套的正确性:asyncpg支持嵌套事务,顶级事务回滚后,所有子事务的变更都会被撤销,完美适配测试场景。
内容的提问来源于stack exchange,提问作者Ahmed Shawky Ahmed
相关产品推荐
相关产品推荐

