FastAPI StreamingResponse未通过生成器函数实现流式输出问题
问题:FastAPI StreamingResponse无法流式返回ChatGPT响应,而是一次性返回全部内容
我开发了一个简易FastAPI应用,接收查询请求后流式返回ChatGPT API的响应。ChatGPT已实现流式返回,结果可实时打印到控制台,但FastAPI的StreamingResponse未实现流式输出,而是一次性返回全部内容,无法定位问题原因。
FastAPI应用代码
import os import time import openai import fastapi from fastapi import Depends, HTTPException, status, Request from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials from fastapi.responses import StreamingResponse auth_scheme = HTTPBearer() app = fastapi.FastAPI() openai.api_key = os.environ["OPENAI_API_KEY"] def ask_statesman(query: str): #prompt = router(query) completion_reason = None response = "" while not completion_reason or completion_reason == "length": openai_stream = openai.ChatCompletion.create( model="gpt-3.5-turbo", messages=[{"role": "user", "content": query}], temperature=0.0, stream=True, ) for line in openai_stream: completion_reason = line["choices"][0]["finish_reason"] if "content" in line["choices"][0].delta: current_response = line["choices"][0].delta.content print(current_response) yield current_response time.sleep(0.25) @app.post("/") async def request_handler(auth_key: str, query: str): if auth_key != "123": raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, detail="Invalid authentication credentials", headers={"WWW-Authenticate": auth_scheme.scheme_name}, ) else: stream_response = ask_statesman(query) return StreamingResponse(stream_response, media_type="text/plain") if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000, debug=True, log_level="debug")
测试代码(test.py)
import requests query = "How tall is the Eiffel tower?" url = "http://localhost:8000" params = {"auth_key": "123", "query": query} response = requests.post(url, params=params, stream=True) for chunk in response.iter_lines(): if chunk: print(chunk.decode("utf-8"))
问题原因及解决方案
核心问题
- 同步生成器与异步路由冲突:
ask_statesman是同步生成器,但路由函数request_handler是异步函数。FastAPI在异步路由中处理同步生成器时,会默认先收集完所有生成器内容再一次性返回,导致流式失效。 - 测试代码读取方式不合适:
iter_lines()需要等待换行符才返回内容,但ChatGPT的流式响应是连续文本块,没有换行,导致客户端无法实时接收内容。
修复方案
方案1:将路由改为同步函数
直接移除路由函数的async关键字,FastAPI会自动在线程池中处理同步逻辑,同步生成器的流式输出即可正常工作:
@app.post("/") def request_handler(auth_key: str, query: str): if auth_key != "123": raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, detail="Invalid authentication credentials", headers={"WWW-Authenticate": auth_scheme.scheme_name}, ) else: stream_response = ask_statesman(query) return StreamingResponse(stream_response, media_type="text/plain")
方案2:将同步生成器转为异步生成器
若需保持路由为异步,可用asyncio.to_thread包装同步生成器,确保事件循环能实时推送内容:
import asyncio async def ask_statesman_async(query: str): def sync_gen(): completion_reason = None while not completion_reason or completion_reason == "length": openai_stream = openai.ChatCompletion.create( model="gpt-3.5-turbo", messages=[{"role": "user", "content": query}], temperature=0.0, stream=True, ) for line in openai_stream: completion_reason = line["choices"][0]["finish_reason"] if "content" in line["choices"][0].delta: current_response = line["choices"][0].delta.content print(current_response) yield current_response time.sleep(0.25) gen = sync_gen() while True: try: yield next(gen) await asyncio.sleep(0) # 让出事件循环,保证响应实时推送 except StopIteration: break @app.post("/") async def request_handler(auth_key: str, query: str): if auth_key != "123": raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, detail="Invalid authentication credentials", headers={"WWW-Authenticate": auth_scheme.scheme_name}, ) else: return StreamingResponse(ask_statesman_async(query), media_type="text/plain")
修复测试代码
将iter_lines()改为iter_content(),并添加flush=True确保实时打印:
import requests query = "How tall is the Eiffel tower?" url = "http://localhost:8000" params = {"auth_key": "123", "query": query} response = requests.post(url, params=params, stream=True) for chunk in response.iter_content(chunk_size=1024): if chunk: print(chunk.decode("utf-8"), end="", flush=True)
内容的提问来源于stack exchange,提问作者Robert Ritz
相关产品推荐
相关产品推荐

