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

如何在FastAPI中结合anyio.TaskGroup与StreamingResponse实现结果实时流式输出?

解决方案:在FastAPI中用AnyIO任务组实现流式并发响应

你的需求非常典型——用结构化并发处理多个外部请求,同时实时流式返回结果,这确实会碰到任务组与异步生成器的冲突问题,因为AnyIO(和Trio)的任务组要求所有子任务在上下文退出前完成,而yield会暂停生成器,打破这个生命周期约束。下面给你几个可行的方案:

方案1:用AnyIO内存对象流+后台任务组实现异步生成器

核心思路是把任务组的运行放在后台任务里,同时用异步生成器持续从内存对象流中读取结果,这样就能在任务组运行的同时流式返回数据,而不用等所有任务完成。

修改后的代码如下:

import anyio
from fastapi import FastAPI
from fastapi.responses import StreamingResponse

app = FastAPI()

@app.get("/")
async def root():
    # 创建异步生成器并传入StreamingResponse
    return StreamingResponse(stream_results(), media_type="text/plain")

async def stream_results():
    send_stream, receive_stream = anyio.create_memory_object_stream(max_buffer_size=5)
    
    async def run_tasks():
        async with anyio.create_task_group() as tg:
            async with send_stream:
                for num in range(5):
                    tg.start_soon(sometask, num, send_stream.clone())
    
    # 启动任务组的后台任务
    async with anyio.create_task_group() as bg_tg:
        bg_tg.start_soon(run_tasks)
        
        # 从接收流读取数据并yield,实现流式输出
        async with receive_stream:
            async for entry in receive_stream:
                yield entry

async def sometask(num, send_stream):
    await anyio.sleep(1)
    async with send_stream:
        await send_stream.send(f'number {num}\n')

if __name__ == "__main__":
    import uvicorn
    uvicorn.run(app)

关键改动说明:

  • 把任务组的执行放到run_tasks函数中,作为后台任务启动,这样任务组的生命周期和生成器的读取操作解耦
  • stream_results作为异步生成器,在后台任务运行的同时,持续从receive_stream读取结果并yield,直接被StreamingResponse消费
  • 设置max_buffer_size避免内存溢出(如果任务返回结果很快的话)

方案2:使用trio_util.trio_async_generator(仅限Trio后端)

如果你用FastAPI的Trio后端(启动时指定uvicorn.run(app, loop="trio")),是可以直接使用trio_util.trio_async_generator的。这个装饰器会帮你处理任务组和生成器的生命周期冲突,允许你在任务组上下文中yield结果。

示例代码大概是这样:

from trio_util import trio_async_generator
# ... 其他导入保持不变

@app.get("/")
async def root():
    return StreamingResponse(main(), media_type="text/plain")

@trio_async_generator
async def main():
    async with anyio.create_task_group() as tg:
        for num in range(5):
            # 这里直接把生成器对象传给任务,任务完成后yield结果
            tg.start_soon(sometask_with_yield, num, yield)

async def sometask_with_yield(num, yield_func):
    await anyio.sleep(1)
    await yield_func(f'number {num}\n')

不过要注意:

  • 这个方案依赖trio_util库,需要额外安装
  • 只能在Trio后端下工作,如果你的FastAPI用的是默认的asyncio后端,这个方法不适用

为什么原来的代码会先聚合所有结果?

你原来的代码里,main函数是先等任务组完全结束(async with anyio.create_task_group()上下文退出),才开始读取receive_stream的所有内容并放到列表里,最后一次性返回,自然没办法实现流式输出。而上面的方案都是让结果的读取和任务的执行同时进行,从而实现实时流式返回。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.01 02:17:34