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

Django集成LangChain+OpenAI实现SSE流式响应失败求助

Django + LangChain + OpenAI SSE流式输出问题排查与修复

核心问题分析

  • 回调处理无效:使用StreamingStdOutCallbackHandler仅将内容打印到终端,未将流式输出传递给SSE响应生成器,客户端无法实时接收数据。
  • 循环逻辑错误:QA查询代码位于for循环外部,仅会处理最后一个问题,且是同步获取完整结果,未触发流式输出机制。
  • 未拆分流式片段:直接返回完整响应结果,没有按OpenAI的流式输出拆分单个token并按SSE格式返回。

修复方案

1. 自定义流式回调Handler

创建专属回调类,捕获LangChain输出的每个token,并传递给SSE生成器。

2. 修复循环逻辑

将QA查询逻辑移入循环内部,确保每个问题都触发流式输出。

3. 配置SSE响应头部

添加禁用缓存的头部,避免中间件或客户端缓存流数据,保证实时性。

修改后的完整代码

from langchain.callbacks.base import BaseCallbackHandler
from rest_framework.decorators import api_view
from rest_framework.response import Response
from django.http import StreamingHttpResponse
from queue import Queue
import threading

# 自定义流式回调Handler,捕获每一块输出内容
class SSEStreamingCallbackHandler(BaseCallbackHandler):
    def __init__(self, queue):
        self.queue = queue

    def on_llm_new_token(self, token: str, **kwargs) -> None:
        # 将每个token转义换行符后按SSE格式存入队列
        self.queue.put(f"data: {token.replace('\n', '\\n')}\n\n")

@api_view(['GET','POST'])
def sse_view(request):
    if request.method != 'POST':
        return Response({"message": "Use the POST method for the response"})

    url = request.data.get("url")
    questions = request.data.get("questions")
    prompt = request.data.get("promptName")

    if not url or not questions:
        return Response({"message": "Please provide valid URL and questions"})

    # Process the documents received from the user
    try:
        doc_store = process_url(url)
        if not doc_store:
            return Response({"message": "PDF document not loaded"})
    except Exception as e:
        return Response({"message": "Error loading PDF document"})

    custom_prompt_template = set_custom_prompt(url,prompt)
    # Load and process documents
    loader = DirectoryLoader(DATA_PATH, glob='*.pdf', loader_cls=PyPDFLoader)
    documents = loader.load()   

    text_splitter = CustomRecursiveCharacterTextSplitter(chunk_size=1500, chunk_overlap=30)
    texts = text_splitter.split_documents(documents)

    #Creating Embeddings using OpenAI
    embeddings = OpenAIEmbeddings(chunk_size= 16, openai_api_key= openai_gpt_key,)

    db = FAISS.from_documents(texts, embeddings)
    db.save_local(DB_FAISS_PATH)
    search_kwargs = {
        'k': 30,
        'fetch_k':100,
        'maximal_marginal_relevance': True,
        'distance_metric': 'cos',
    }
    retriever=db.as_retriever(search_kwargs=search_kwargs)

    # get the list of questions from the body
    questionList = request.data['questions']

    #This function is to generate responses from OpenAi
    def openai_response_generator():
        for question in questionList:
            queue = Queue()
            callback = SSEStreamingCallbackHandler(queue=queue)
            
            # Create an instance of ChatOpenAI with custom callback
            llm = ChatOpenAI(
                model_name="gpt-3.5-turbo-16k",
                streaming=True,
                callbacks=[callback],
                temperature=0,
                openai_api_key= openai_gpt_key,
            )

            qa = RetrievalQA.from_chain_type(
                llm=llm,
                chain_type="stuff",
                retriever=retriever,
                return_source_documents=True,
                chain_type_kwargs={"prompt": custom_prompt_template},
            )

            # 启动线程执行QA查询,避免阻塞生成器
            def run_qa():
                qa({'query': question})
                # 查询结束后放入终止标记
                queue.put(None)

            threading.Thread(target=run_qa).start()

            # 从队列中获取token并yield到响应流
            while True:
                token_data = queue.get()
                if token_data is None:
                    break
                yield token_data
            # 每个问题结束后发送分隔标记(可选)
            yield "data: --- 当前问题回答完成 ---\n\n"

    response = StreamingHttpResponse(openai_response_generator(), content_type="text/event-stream")
    # 添加SSE必要头部,禁用缓存
    response['Cache-Control'] = 'no-cache'
    response['Connection'] = 'keep-alive'
    response['X-Accel-Buffering'] = 'no'  # 防止反向代理(如Nginx)缓存流数据
    return response

关键修改说明

  • 自定义回调:通过on_llm_new_token捕获每个输出token,存入队列供生成器读取。
  • 多线程处理:用线程启动QA查询,避免阻塞生成器,确保流式token能实时返回。
  • 头部配置:添加Cache-Control、X-Accel-Buffering等头部,彻底避免缓存导致的延迟。
  • 循环修复:将QA查询逻辑移入循环,每个问题独立触发流式输出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 07:09:58