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

Pytest测试FastAPI接口报错:Future关联不同事件循环

解决FastAPI+Elasticsearch测试中的事件循环不匹配问题

问题场景

在测试FastAPI创建Elasticsearch数据源的接口时,触发了事件循环不匹配的RuntimeError,相关代码及错误信息如下:

测试代码

def make_test_source(
    client: TestClient,
    name: str = "test-source",
    quantize: bool = False,
    embedding_model: EmbeddingModel = EmbeddingModel.TEXT_EMBEDDING_ADA_002,
    **kwargs,
):
    return client.post(
        url="/api/v1/sources",
        json={
            "source": name,
            "contains_secret_data": False,
            "adct_id": "ADCT-12345",
            "embedding_model": embedding_model,
            "quantize": quantize,
            **kwargs,
        },
    )

@pytest.mark.anyio
async def test_create_source(es, client):
    make_test_source(client=client, name="test-source")
    assert await es.indices.exists(index="test-source")
    assert "test-source" in (await get_global_sources(es))

Fixtures定义

@pytest.fixture
def anyio_backend():
    return "asyncio"

@pytest.fixture
def client():
    app = create_app(test=True)
    app.dependency_overrides[get_user] = get_user_override
    with TestClient(app) as c:
        yield c

@pytest.fixture
async def es(client):
    """Initialize Elasticsearch client and database.

    Sets up metadata and permissions indices and removes all indices
    after each test.

    Pass this function to every test case that uses Elasticsearch, even
    if it doesn't call the ES client directly.
    """
    _es = connections["es"] 
    await create_metadata_index(_es)
    await create_permissions_index(_es)
    yield _es
    await clear_all_indices(_es)

接口实现

@router.post(
    path="/",
    summary="Create a new source",
    status_code=HTTP_201_CREATED)
async def create_source(
    body: CreateSourceRequest,
    user: Annotated[User, Depends(required_user)],
    es: ElasticSearchDep,
):
    embed_dims = embedding_models[body.embedding_model].embed_dim
    quantize = body.quantize
    try:
        # --- 错误触发点 ---
        await es.indices.create(
            index=body.source, mappings=get_document_mapping(embed_dims, quantize=quantize)
        ) 
    ...

应用生命周期配置(apps.py)

connections = {} 

@contextlib.asynccontextmanager
async def lifespan_test(app: FastAPI):
    connections["es"] = AsyncElasticsearch(
        hosts=[
            {
                "host": "0.0.0.0",
                "port": 9200,
                "scheme": "https",
            }
        ],
        basic_auth=("elastic", settings.ES_TEST_PASSWORD),
        verify_certs=False,
        ssl_show_warn=False,
    )
    yield
    await connections["es"].close()

def create_app(test: bool = False):
    app = FastAPI(title="test", lifespan=lifespan_test)
    
    app.add_middleware(
        CORSMiddleware,
        allow_origins=["*"],
        allow_methods=["*"],
        allow_headers=["*"],
        expose_headers=["*"],
    )
    app.add_middleware(BaseHTTPMiddleware, dispatch=catch_exceptions_middleware)
    app.include_router(router)

    return app

app = create_app()

错误信息

RuntimeError: Task <Task pending name='starlette.middleware.base.BaseHTTPMiddleware.__call__.<locals>.call_next.<locals>.coro' coro=<BaseHTTPMiddleware.__call__.<locals>.call_next.<locals>.coro() running at /Users/.../venv/lib/python3.11/site-packages/starlette/middleware/base.py:151> cb=[TaskGroup._spawn.<locals>.task_done() at /Users/.../venv/lib/python3.11/site-packages/anyio/_backends/_asyncio.py:661]> got Future <Future pending> attached to a different loop

问题原因

核心矛盾在于事件循环不匹配:

  • 同步定义的client fixture中,TestClient会创建独立的临时事件循环处理请求
  • connections["es"]中的Elasticsearch连接是在lifespan中初始化的,绑定的是pytest-anyio管理的测试主循环
  • 当接口请求调用ES时,请求所在的临时循环与ES连接绑定的主循环不一致,导致跨循环Future传递错误

解决步骤

1. 将Client Fixture改为异步

使用AsyncTestClient(需确保安装starlette>=0.20),并将fixture改为异步定义,确保请求运行在测试主循环中:

@pytest.fixture
async def client():
    app = create_app(test=True)
    app.dependency_overrides[get_user] = get_user_override
    async with AsyncTestClient(app) as c:
        yield c

2. 调整ES连接的初始化逻辑

将ES连接的初始化从lifespan移到es fixture中,避免跨循环绑定:

  • 修改apps.py,拆分ES连接的初始化/关闭逻辑:
connections = {} 

async def init_es_connection():
    connections["es"] = AsyncElasticsearch(
        hosts=[
            {
                "host": "0.0.0.0",
                "port": 9200,
                "scheme": "https",
            }
        ],
        basic_auth=("elastic", settings.ES_TEST_PASSWORD),
        verify_certs=False,
        ssl_show_warn=False,
    )
    return connections["es"]

async def close_es_connection():
    await connections["es"].close()

def create_app(test: bool = False):
    app = FastAPI(title="test")
    
    # 生产环境保留lifespan管理,测试环境交给fixture处理
    if not test:
        @contextlib.asynccontextmanager
        async def lifespan_prod(app: FastAPI):
            await init_es_connection()
            yield
            await close_es_connection()
        app.lifespan = lifespan_prod
    
    # 中间件和路由配置不变
    app.add_middleware(
        CORSMiddleware,
        allow_origins=["*"],
        allow_methods=["*"],
        allow_headers=["*"],
        expose_headers=["*"],
    )
    app.add_middleware(BaseHTTPMiddleware, dispatch=catch_exceptions_middleware)
    app.include_router(router)

    return app
  • 更新es fixture:
@pytest.fixture
async def es():
    _es = await init_es_connection()
    await create_metadata_index(_es)
    await create_permissions_index(_es)
    yield _es
    await clear_all_indices(_es)
    await close_es_connection()

3. 修改测试函数为异步调用

将make_test_source改为异步函数,并在测试中用await调用:

async def make_test_source(
    client: AsyncTestClient,
    name: str = "test-source",
    quantize: bool = False,
    embedding_model: EmbeddingModel = EmbeddingModel.TEXT_EMBEDDING_ADA_002,
    **kwargs,
):
    return await client.post(
        url="/api/v1/sources",
        json={
            "source": name,
            "contains_secret_data": False,
            "adct_id": "ADCT-12345",
            "embedding_model": embedding_model,
            "quantize": quantize,
            **kwargs,
        },
    )

@pytest.mark.anyio
async def test_create_source(es, client):
    await make_test_source(client=client, name="test-source")
    assert await es.indices.exists(index="test-source")
    assert "test-source" in (await get_global_sources(es))

验证

修改完成后,所有异步操作(请求调用、ES操作)都会运行在同一个测试事件循环中,跨循环的Future传递问题将被解决。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 12:25:29