包装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)
关键注意事项
正确调用生成器:必须使用普通
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())处理响应流:
stream.aiter_raw()返回原始响应字节,若后端返回文本/JSON流,可根据需求改用stream.aiter_lines()或stream.aiter_text(),解码为字符串后再放入队列。线程资源管理:将线程设为
daemon=True,主线程退出时会自动终止后台线程,避免资源泄漏。
内容的提问来源于stack exchange,提问作者thomas
相关产品推荐
相关产品推荐

