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

求助: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 11:22:55