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
相关产品推荐
相关产品推荐

