求助:LangChain OpenAI流式响应无法通过FastAPI端点返回
问题排查与解决:FastAPI无法流式返回OpenAI响应
核心问题分析
当前代码存在三个关键问题导致流式失效:
- 同步阻塞调用:
askQuestion中调用self.chain({"question": question})会等待整个响应生成完成才返回结果,此时generate_tokens遍历的是完整字符串的单个字符,并非实时流式token。 - 队列未被正确利用:自定义回调将token放入同步队列,但生成器并未从队列取token,反而使用了同步返回的完整结果。
- LangChain配置过时:使用了已弃用的
callback_manager,且未采用异步适配的回调和chain调用方式,无法适配FastAPI的异步环境。
修正方案与代码调整
1. 替换为异步回调与异步队列
使用异步队列和异步回调处理器,适配FastAPI的异步事件循环:
import asyncio import sys from langchain.callbacks.base import AsyncCallbackHandler # 异步队列存储流式token token_queue = asyncio.Queue() stop_signal = "###STREAM_END###" class AsyncStreamingHandler(AsyncCallbackHandler): async def on_llm_start(self, serialized: dict[str, any], prompts: list[str], **kwargs: any) -> None: # 启动前清空队列 while not token_queue.empty(): try: token_queue.get_nowait() except asyncio.QueueEmpty: break async def on_llm_new_token(self, token: str, **kwargs: any) -> None: # 将新token放入异步队列 await token_queue.put(token) sys.stdout.write(token) sys.stdout.flush() async def on_llm_end(self, response, **kwargs: any) -> None: # 发送流式结束信号 await token_queue.put(stop_signal)
2. 重构askQuestion为异步方法
修改方法为异步类型,启动chain的异步任务而非同步等待结果:
async def askQuestion(self, collection_id, question): collection_name = "collection-" + str(collection_id) # 配置异步流式LLM,使用新版callbacks参数 self.llm = ChatOpenAI( model_name=self.model_name, temperature=self.temperature, openai_api_key=os.environ.get('OPENAI_API_KEY'), streaming=True, verbose=VERBOSE, callbacks=[AsyncStreamingHandler()] ) self.memory = ConversationBufferMemory( memory_key="chat_history", return_messages=True, output_key='answer' ) chroma_Vectorstore = Chroma( collection_name=collection_name, embedding_function=self.embeddingsOpenAi, client=self.chroma_client ) self.chain = ConversationalRetrievalChain.from_llm( self.llm, chroma_Vectorstore.as_retriever(similarity_search_with_score=True), return_source_documents=True, verbose=VERBOSE, memory=self.memory ) # 异步启动chain任务,不阻塞等待结果 await asyncio.create_task(self.chain.arun(question=question))
3. 调整FastAPI路由的流式生成器
改为异步生成器,从队列实时获取token并返回:
from fastapi import StreamingResponse, HTTPException @router.post("/collection/{collection_id}/ask_question") async def ask_question(collection_id: str, request: Request): try: form_data = await request.form() question = form_data["question"] async def generate_tokens(): # 触发提问任务 await thread_handler.askQuestion(collection_id, question) # 循环从队列取token,直到收到结束信号 while True: token = await token_queue.get() if token == stop_signal: break # 以utf-8编码返回token,media_type匹配纯文本流式 yield token.encode('utf-8') # 使用text/plain作为媒体类型,适合流式文本输出 return StreamingResponse(generate_tokens(), media_type="text/plain") except requests.exceptions.ConnectionError as e: raise HTTPException(status_code=500, detail="连接服务器失败") except Exception as e: raise HTTPException(status_code=500, detail=str(e))
关键修改说明
- 异步适配:全部核心逻辑改为异步,避免阻塞FastAPI的事件循环,确保token能实时推送。
- 队列利用:通过异步队列在回调和生成器之间传递token,实现真正的流式传输。
- LangChain版本兼容:弃用旧版
callback_manager,改用新版callbacks参数,确保流式功能正常触发。 - 媒体类型调整:将
media_type改为text/plain,更适合纯文本流式输出;若需JSON格式,可将yield内容改为yield f'"{token}"\n'并对应调整media_type为application/json。
内容的提问来源于stack exchange,提问作者devZ
相关产品推荐
相关产品推荐

