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

如何用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通知逻辑等

关键注意事项

  1. 数据库操作必须正确释放连接:确保db_handler中的方法使用async with管理连接,比如:
    async def fetch_todos():
        async with await get_connection() as conn:
            return await conn.fetch("SELECT * FROM todos")
    
  2. 避免遗漏await:原代码中如果db_handler.insert或fetch_todos是异步方法,必须添加await,否则会返回协程对象导致逻辑错误。
  3. 事务嵌套的正确性:asyncpg支持嵌套事务,顶级事务回滚后,所有子事务的变更都会被撤销,完美适配测试场景。

内容的提问来源于stack exchange,提问作者Ahmed Shawky Ahmed

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 10:10:31