Skip to content
🔗 分享本题
查看我的学习进度 →

LLM Token 流经后端异步迭代器和 SSE 到浏览器增量渲染的端到端方案

🧠 图解记忆: 流式输出是端到端管道,后端逐块转发、前端增量渲染,并处理取消、断线和背压。

💡 答案要点

Streaming = LLM边生成边返回token,首屏时间从5s→0.3s

为什么需要流式输出

非流式:用户等5秒 → 看到完整回答(体验差)
流式:  0.3秒开始显示第一个字 → 逐字显示 → 5秒显示完(体验好)

指标对比:
TTFT(首Token时间):5000ms → 300ms(-94%)
用户感知等待时间:5s → 0.3s(-94%)
用户满意度:↑ 显著提升

后端实现(FastAPI SSE)

展开 Python 代码示例(64 行)
python
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from openai import OpenAI
import asyncio
import json

app = FastAPI()
client = OpenAI()

@app.post("/chat/stream")
async def chat_stream(request: ChatRequest):
    """流式聊天接口"""

    async def generate():
        try:
            # 开启流式模式
            stream = client.chat.completions.create(
                model="gpt-4o-mini",
                messages=[{"role": "user", "content": request.message}],
                stream=True,          # 关键参数
                max_tokens=2048,
            )

            # 逐chunk发送
            for chunk in stream:
                delta = chunk.choices[0].delta

                if delta.content:
                    # SSE格式:data: {json}\n\n
                    data = json.dumps({
                        "type": "content",
                        "content": delta.content
                    }, ensure_ascii=False)
                    yield f"data: {data}\n\n"

                # 工具调用流式处理
                if delta.tool_calls:
                    for tc in delta.tool_calls:
                        data = json.dumps({
                            "type": "tool_call",
                            "tool": tc.function.name if tc.function.name else "",
                            "args_chunk": tc.function.arguments if tc.function.arguments else ""
                        })
                        yield f"data: {data}\n\n"

                # 让出事件循环,避免阻塞
                await asyncio.sleep(0)

            # 发送结束信号
            yield f"data: {json.dumps({'type': 'done'})}\n\n"

        except Exception as e:
            error_data = json.dumps({"type": "error", "message": str(e)})
            yield f"data: {error_data}\n\n"

    return StreamingResponse(
        generate(),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "X-Accel-Buffering": "no",   # 禁止Nginx缓冲
            "Connection": "keep-alive",
        }
    )

前端实现(EventSource / fetch)

展开 Javascript 代码示例(54 行)
javascript
// 方式1:EventSource(简单,仅支持GET)
const es = new EventSource('/chat/stream?message=你好');

es.onmessage = (event) => {
    const data = JSON.parse(event.data);

    if (data.type === 'content') {
        // 追加到显示区域
        document.getElementById('output').textContent += data.content;
    } else if (data.type === 'done') {
        es.close();
    } else if (data.type === 'error') {
        console.error(data.message);
        es.close();
    }
};

// 方式2:fetch + ReadableStream(支持POST,更灵活)
async function streamChat(message) {
    const response = await fetch('/chat/stream', {
        method: 'POST',
        headers: {'Content-Type': 'application/json'},
        body: JSON.stringify({message})
    });

    const reader = response.body.getReader();
    const decoder = new TextDecoder();
    let buffer = '';

    while (true) {
        const {done, value} = await reader.read();
        if (done) break;

        buffer += decoder.decode(value, {stream: true});

        // 按SSE格式解析
        const lines = buffer.split('\n');
        buffer = lines.pop();  // 保留不完整的行

        for (const line of lines) {
            if (line.startsWith('data: ')) {
                const jsonStr = line.slice(6);
                if (jsonStr === '[DONE]') return;

                try {
                    const data = JSON.parse(jsonStr);
                    if (data.type === 'content') {
                        appendToOutput(data.content);
                    }
                } catch (e) {}
            }
        }
    }
}

中间件处理(LangChain流式)

展开 Python 代码示例(43 行)
python
from langchain.callbacks.streaming_stdout import StreamingStdOutCallbackHandler
from langchain.callbacks.base import BaseCallbackHandler

class CustomStreamHandler(BaseCallbackHandler):
    """自定义流式回调"""

    def __init__(self, queue: asyncio.Queue):
        self.queue = queue

    def on_llm_new_token(self, token: str, **kwargs):
        """每生成一个token触发"""
        self.queue.put_nowait(token)

    def on_llm_end(self, response, **kwargs):
        """生成结束"""
        self.queue.put_nowait(None)  # 发送结束信号

# 在FastAPI中使用
@app.post("/langchain/stream")
async def langchain_stream(request: ChatRequest):
    queue = asyncio.Queue()
    handler = CustomStreamHandler(queue)

    async def generate():
        # 在后台线程运行LangChain(避免阻塞事件循环)
        import threading

        def run_chain():
            llm = ChatOpenAI(streaming=True, callbacks=[handler])
            llm.invoke(request.message)

        thread = threading.Thread(target=run_chain)
        thread.start()

        # 从队列读取token并发送
        while True:
            token = await queue.get()
            if token is None:
                yield f"data: {json.dumps({'type': 'done'})}\n\n"
                break
            yield f"data: {json.dumps({'type': 'content', 'content': token})}\n\n"

    return StreamingResponse(generate(), media_type="text/event-stream")

生产注意事项

python
# 1. Nginx配置(防止缓冲)
nginx_config = """
location /chat/stream {
    proxy_pass http://backend;
    proxy_buffering off;           # 关键:禁用缓冲
    proxy_cache off;
    proxy_set_header Connection '';
    proxy_http_version 1.1;
    chunked_transfer_encoding on;
}
"""

# 2. 超时设置(流式响应时间长)
# 普通接口:30s超时
# 流式接口:600s超时(10分钟)

# 3. 断点续传(网络中断后继续)
class StreamResumer:
    def __init__(self):
        self.cache = {}  # session_id → 已生成内容

    def resume_stream(self, session_id: str, offset: int):
        """从offset位置继续输出"""
        cached = self.cache.get(session_id, "")
        if len(cached) > offset:
            # 先发送已缓存的部分
            return cached[offset:], True
        return "", False

面试话术:

"流式输出核心是Server-Sent Events(SSE)协议:后端yield逐个token,前端EventSource或fetch ReadableStream实时接收追加。TTFT从5s降到300ms,用户体验大幅提升。生产上3个关键点:1)Nginx必须关闭proxy_buffering,否则还是等全部生成才返回;2)流式接口超时设置要长,普通30s会被截断;3)用asyncio.sleep(0)让出事件循环,避免阻塞其他请求。工具调用也可以流式传递,让用户看到Agent正在思考的过程。"

📚 参考:MDN:Server-Sent Events(流式输出标准)