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

包装FastAPI流式响应为同步生成器时遇Task got bad yield错误

解决RuntimeError: Task got bad yield问题

问题根源

RuntimeError: Task got bad yield错误的核心原因是普通生成器被错误地放入asyncio Task中执行——asyncio的Task仅能处理使用await的协程,普通生成器的yield语法会被Task判定为非法操作。

另外你的代码存在可优化点:带timeout的queue.get()会频繁触发queue.Empty异常,属于没必要的额外循环开销。

修正后的代码

import queue
import httpx
from threading import Thread
import asyncio

# 替换为你的FastAPI后端地址
BACKEND = "http://your-fastapi-backend"

def process_stream(input_str):
    q = queue.Queue()
    job_done = object()

    def task():
        # 为线程创建独立的事件循环
        loop = asyncio.new_event_loop()
        asyncio.set_event_loop(loop)
        try:
            loop.run_until_complete(process_stream_async(input_str, q))
        finally:
            loop.close()
            q.put(job_done)

    Thread(target=task, daemon=True).start()

    while True:
        next_token = q.get()  # 阻塞等待,无需timeout
        if next_token is job_done:
            break
        # 若需要字符串而非字节,可添加解码:next_token.decode('utf-8')
        yield next_token

async def process_stream_async(input_str, q):
    payload = {"input": input_str}
    async with httpx.AsyncClient(timeout=None) as client:
        async with client.stream(
            "POST",
            url=f"{BACKEND}/process-stream",
            json=payload
        ) as stream:
            # 如果后端返回JSON流(如每行一个JSON对象),可改用stream.aiter_lines()
            async for item in stream.aiter_raw():
                q.put(item)

关键注意事项

  1. 正确调用生成器:必须使用普通for循环调用,绝对不能在async函数中用await或async for执行该生成器:

    # 正确用法
    for token in process_stream("some input"):
        print(token)
    

    若需在async环境中使用,需将普通生成器包装为异步迭代器:

    import asyncio
    
    async def async_process_stream(input_str):
        for token in process_stream(input_str):
            yield token
    
    # 在async函数中使用
    async def main():
        async for token in async_process_stream("some input"):
            print(token)
    
    asyncio.run(main())
    
  2. 处理响应流:stream.aiter_raw()返回原始响应字节,若后端返回文本/JSON流,可根据需求改用stream.aiter_lines()或stream.aiter_text(),解码为字符串后再放入队列。

  3. 线程资源管理:将线程设为daemon=True,主线程退出时会自动终止后台线程,避免资源泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 18:02:42