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

基于LangChain与FastAPI的RAG聊天机器人流式响应如何正确携带元数据

基于LangChain与FastAPI的RAG聊天机器人流式响应如何正确携带元数据

这个问题我之前做RAG流式响应时也碰到过——直接把元数据和文本内容一起yield,很容易因为浏览器的SSE chunk合并策略导致两者混在一起,加延迟完全是治标不治本,网络波动或者浏览器的缓冲逻辑随时可能让问题复发。给你几个更靠谱的方案,都是适配你当前LangChain+FastAPI技术栈的:

方案一:用SSE标准事件类型区分内容与元数据(最推荐)

SSE(Server-Sent Events)本身就支持自定义事件类型,我们可以把LLM的流式内容和元数据分成两个不同的事件发送,前端只需要监听对应的事件类型就能分别处理,完全不会混在一起。

后端代码修改

调整你的arag_stream_pinecone函数里的yield逻辑,把普通内容和元数据包装成符合SSE规范的事件:

async for chunk in self.streaming_llm.astream(prompt_text):
    if chunk:
        # 发送"content"类型的事件,携带LLM回复片段
        yield f"event: content\ndata: {chunk.content or ''}\n\n"
# 所有内容发送完成后,发送"metadata"类型的事件
if metadata:
    # 如果元数据是复杂结构,可以用json.dumps转成字符串
    # import json
    # yield f"event: metadata\ndata: {json.dumps({'ids': idList, 'metadata': metadata})}\n\n"
    yield f"event: metadata\ndata: {metadata}\n\n"

注意每个SSE事件必须以\n\n结尾,浏览器才能正确识别为一个独立事件。

前端处理示例(JavaScript)

前端通过监听不同的事件类型,分别处理内容和元数据:

const eventSource = new EventSource('/message/private?stream=true&document_type=your_type');
let fullResponse = '';

// 监听LLM内容片段
eventSource.addEventListener('content', (event) => {
    fullResponse += event.data;
    // 更新UI显示实时回复
    document.getElementById('chat-response').textContent = fullResponse;
});

// 监听元数据事件
eventSource.addEventListener('metadata', (event) => {
    const metadata = event.data;
    // 这里调用你的元数据查询接口
    fetch('/get-metadata', {
        method: 'POST',
        body: JSON.stringify({ids: metadata.split(',')})
    }).then(res => res.json()).then(metaDataDetails => {
        // 处理元数据详情,比如展示引用来源
    });
    // 可选:元数据接收完成后关闭SSE连接
    eventSource.close();
});

// 错误处理
eventSource.addEventListener('error', (err) => {
    console.error('流式响应出错:', err);
    eventSource.close();
});

这个方案完全符合SSE规范,是最标准的做法,没有兼容性问题,也不需要担心内容和元数据混在一起。

方案二:用特殊分隔标记包裹元数据(快速过渡方案)

如果不想大改SSE事件结构,可以给元数据加一个唯一的特殊分隔符,前端在接收所有chunk后,把分隔符前后的内容拆分,分别作为回复内容和元数据处理。

后端代码修改

async for chunk in self.streaming_llm.astream(prompt_text):
    if chunk:
        yield chunk.content or ''
# 用一个极不可能出现在LLM回复里的分隔符包裹元数据
if metadata:
    yield f"\n===__METADATA_MARKER__==={metadata}===__METADATA_MARKER__==="

分隔符要选得足够特殊,比如包含大小写、下划线、特殊符号,避免和LLM的正常回复冲突,也可以在prompt里加一句“不要在回复中包含===METADATA_MARKER===这类特殊标记”。

前端处理示例

const eventSource = new EventSource('/message/private?stream=true&document_type=your_type');
let fullResponse = '';
const metadataMarker = '===__METADATA_MARKER__===';

eventSource.onmessage = (event) => {
    if (event.data.includes(metadataMarker)) {
        // 拆分出内容和元数据
        const [contentPart, metaPart] = event.data.split(metadataMarker);
        fullResponse += contentPart;
        const metadata = metaPart.trim();
        // 处理元数据
        eventSource.close();
    } else {
        fullResponse += event.data;
        // 更新UI
    }
};

这个方案改动最小,适合快速调整,但不如SSE事件类型方案规范。

方案三:改用WebSocket实现更灵活的流式通信

如果你的场景允许调整通信协议,WebSocket比SSE更灵活,可以直接发送结构化的JSON消息,每个消息带一个type字段区分内容和元数据,前端处理起来更直观。

后端WebSocket代码示例

from fastapi import WebSocket, WebSocketDisconnect
import json

@router.websocket("/ws/chat/private")
async def chat_websocket(
    websocket: WebSocket,
    document_type: str,
    current_user: dict = Depends(auth_service.get_current_user)
):
    await websocket.accept()
    try:
        # 接收前端发送的提问和聊天历史
        client_data = await websocket.receive_json()
        message = client_data["message"]
        chat_history = client_data["chat_history"]
        
        # 执行RAG流式逻辑
        rephrased_query = self.rephrase_query(message, chat_history)
        # ... 省略向量查询、文档检索的逻辑,和你原来的代码一致 ...
        metadata = ','.join(idList) if authorization != "public" else ""

        # 流式发送LLM回复
        async for chunk in self.streaming_llm.astream(prompt_text):
            if chunk:
                await websocket.send_json({
                    "type": "content",
                    "data": chunk.content or ''
                })
        # 发送元数据
        await websocket.send_json({
            "type": "metadata",
            "data": metadata
        })
    except WebSocketDisconnect:
        logging.info("WebSocket连接已断开")
    except Exception as e:
        logging.error(f"WebSocket错误: {e}")
        await websocket.send_json({"type": "error", "data": str(e)})

前端WebSocket处理示例

const socket = new WebSocket(`ws://localhost:8000/ws/chat/private?document_type=your_type`);
let fullResponse = '';

socket.onopen = () => {
    // 发送提问和聊天历史
    socket.send(JSON.stringify({
        message: userInput,
        chat_history: chatHistory
    }));
};

socket.onmessage = (event) => {
    const data = JSON.parse(event.data);
    if (data.type === 'content') {
        fullResponse += data.data;
        // 更新UI
    } else if (data.type === 'metadata') {
        // 处理元数据
        socket.close();
    } else if (data.type === 'error') {
        console.error('错误:', data.data);
    }
};

WebSocket适合需要双向通信的场景(比如前端中途中断请求),但需要调整前端和后端的通信逻辑,比SSE的改动大一些。

总结

  • 优先选方案一(结构化SSE事件):符合标准,改动适中,前端后端处理都清晰,没有兼容性问题。
  • 快速过渡用方案二(分隔符标记):改动最小,但要注意分隔符的唯一性。
  • 复杂双向场景用方案三(WebSocket):灵活性最高,但需要调整通信协议。

完全不推荐加延迟的方案,因为浏览器的chunk合并策略是不可控的,延迟时间短了还是会混,延迟长了影响用户体验。

备注:内容来源于stack exchange,提问作者Abdullah Muhammad Moosa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 11:59:29