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
问题原因
核心矛盾在于事件循环不匹配:
- 同步定义的
clientfixture中,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
- 更新
esfixture:
@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
相关产品推荐
相关产品推荐

