AI流式处理折腾手记

「你这个是不是在用人工识别?」

做实时语音助手时,第一版还是 HTTP 一问一答:用户说完话,等 2–3 秒才看到识别和回复,体验差到被当面吐槽。后来改流式,首包压到 500ms 以内,token 传输、缓冲区、WebSocket 每个环节都有具体坑要对。

为什么写这篇文章

最近在做一个实时语音助手项目,用户对着麦克风说话,系统需要实时识别语音并给出反馈。最初我们用了传统的 HTTP 请求模式,用户说一句话就要等 2-3 秒才能看到结果,体验很差。

用户直接吐槽:“你这个是不是在用人工识别?”

确实,体验太差了。于是开始研究流式处理,把延迟从 2-3 秒降到了 500ms 以内。但这个过程踩了不少坑,特别是 token 传输、缓冲区管理、WebSocket 优化这些细节。把这些经验整理出来,希望能帮到类似需求的同学。

背景和需求

遇到的问题

先说下场景:我们的智能客服系统,用户在网页上输入问题,后端调用 LLM 生成回答。最开始的实现逻辑很简单:

sequenceDiagram participant User as 用户 participant Frontend as 前端 participant Backend as 后端 participant LLM as AI模型 User->>Frontend: 输入问题 Frontend->>Backend: HTTP POST /chat Backend->>LLM: 生成完整回答 Note over Backend,LLM: 等待 2-3 秒 LLM->>Backend: 返回完整回答 Backend->>Frontend: 返回完整回答 Frontend->>User: 显示回答

这个实现的问题很明显:

  1. 用户体验差:用户盯着白屏等 2-3 秒,感觉很卡
  2. 服务端压力大:高并发时大量请求堆积
  3. 资源浪费:生成 100 个 token,等最后一个才返回,其实中间的 token 早就可以展示了

需要达到的效果

  • 首字延迟(TTFT, Time to First Token):< 500ms
  • token 间隔:每 50-100ms 输出一个 token
  • 整体延迟:用户感觉像在"实时对话"

换句话说,就像 GPT 那样的打字机效果,而不是等了半天突然冒出一段话。

实现方案

1. 理解流式处理的基本概念

先说几个关键概念,用大白话解释:

  • 流式处理:不是等生成完整结果再返回,而是生成一点就返回一点,像水流一样源源不断
  • token:AI 处理文本的基本单位,可以理解为一个词或字的一部分
  • TTFT(Time to First Token):从请求发出到收到第一个 token 的时间,这个是用户感知到的"启动延迟"
  • TPOT(Time Per Output Token):每个 token 之间的间隔时间,决定了"打字速度"

打个比方:

  • 传统处理:你在餐厅点餐,厨房把所有菜都做完了一起端上来
  • 流式处理:厨师炒好一道菜就端上来一道菜,你边吃边等

2. 技术选型

后端:WebSocket + Server-Sent Events

对比了几种方案:

方案优点缺点适用场景
HTTP 轮询简单,兼容性好延迟高,资源浪费不推荐用于实时场景
SSE简单,自动重连单向通信适合服务端推送
WebSocket双向通信,实时性强实现复杂,需要额外处理适合双向实时交互

我们选了 WebSocket,因为:

  1. 支持双向通信(用户可以随时打断)
  2. 连接复用,减少握手开销
  3. 各语言都有成熟的库

AI 模型 API:OpenAI Stream API

如果用 OpenAI 的 API,直接用他们的流式接口:

from openai import OpenAI

client = OpenAI()

stream = client.chat.completions.create(
    model="gpt-4",
    messages=[{"role": "user", "content": "写一首诗"}],
    stream=True  # 关键:开启流式
)

for chunk in stream:
    if chunk.choices[0].delta.content:
        print(chunk.choices[0].delta.content, end="")

如果是本地部署模型(比如 vLLM、Ollama),也都有对应的流式 API。

3. 架构设计

整体架构:

graph LR A[用户输入] --> B[前端 WebSocket] B --> C[后端服务] C --> D[流式 AI 模型] D --> C C --> E[WebSocket 推送] E --> F[前端实时展示] style D fill:#f9f,stroke:#333,stroke-width:2px style F fill:#9f9,stroke:#333,stroke-width:2px

前端通过 WebSocket 连接后端,后端调用流式 AI 模型,拿到 token 后立即通过 WebSocket 推回前端,前端实时展示。

4. 后端实现

这里以 Python + FastAPI 为例:

from fastapi import FastAPI, WebSocket
from fastapi.middleware.cors import CORSMiddleware
import json
import asyncio
from openai import AsyncOpenAI

app = FastAPI()

app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_methods=["*"],
    allow_headers=["*"],
)

client = AsyncOpenAI()

@app.websocket("/ws/chat")
async def chat_websocket(websocket: WebSocket):
    await websocket.accept()

    try:
        # 接收用户消息
        data = await websocket.receive_text()
        message = json.loads(data)
        user_input = message.get("content", "")

        # 调用流式 API
        stream = await client.chat.completions.create(
            model="gpt-4",
            messages=[{"role": "user", "content": user_input}],
            stream=True,
            stream_options={"include_usage": True}
        )

        # 逐个 token 推送
        async for chunk in stream:
            if chunk.choices[0].delta.content:
                token = chunk.choices[0].delta.content

                # 发送 token
                await websocket.send_json({
                    "type": "token",
                    "content": token
                })

        # 发送完成信号
        await websocket.send_json({
            "type": "done"
        })

    except Exception as e:
        await websocket.send_json({
            "type": "error",
            "message": str(e)
        })
    finally:
        await websocket.close()

关键点:

  1. AsyncOpenAI 异步调用,避免阻塞
  2. 每收到一个 token 就立即推送,不等待
  3. type 字段区分不同消息类型(token、完成、错误)

5. 前端实现

const ws = new WebSocket('ws://localhost:8000/ws/chat');

ws.onopen = () => {
    console.log('WebSocket 连接已建立');
    // 发送用户消息
    ws.send(JSON.stringify({
        content: '你好,请介绍一下自己'
    }));
};

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

    if (data.type === 'token') {
        // 实时显示 token
        appendToChat(data.content);
    } else if (data.type === 'done') {
        console.log('回答完成');
    } else if (data.type === 'error') {
        console.error('错误:', data.message);
    }
};

function appendToChat(text) {
    const chatContainer = document.getElementById('chat');
    chatContainer.textContent += text;
}

优化点:

  1. 连接保活:定时发送 ping/pong
  2. 断线重连:自动重连机制
  3. 消息队列:防止消息乱序(虽然 WebSocket 一般有序)

踩坑经历

坑 1:WebSocket 连接频繁断开

刚开始实现时,WebSocket 连接经常莫名其妙断开,用户体验很差。

排查过程

  1. 加了日志,发现是服务端主动断开
  2. 检查 Nginx 配置,发现 proxy_read_timeout 是 60 秒
  3. 我们的 AI 回复有时候会超过 60 秒

解决方案

location /ws/ {
    proxy_pass http://backend;
    proxy_http_version 1.1;
    proxy_set_header Upgrade $http_upgrade;
    proxy_set_header Connection "upgrade";
    proxy_read_timeout 300s;  # 改成 5 分钟
    proxy_send_timeout 300s;
}

另外,前端加了心跳机制:

let heartbeatInterval;

function startHeartbeat() {
    heartbeatInterval = setInterval(() => {
        if (ws.readyState === WebSocket.OPEN) {
            ws.send(JSON.stringify({ type: 'ping' }));
        }
    }, 30000);  // 30 秒一次
}

ws.onclose = () => {
    clearInterval(heartbeatInterval);
    // 自动重连
    setTimeout(connect, 1000);
};

坑 2:token 显示速度不均匀

有时候 token 来得很快,有时候又很慢,打字机效果忽快忽慢,很影响体验。

原因分析

  1. AI 模型生成 token 的速度本身就不均匀
  2. 网络抖动也会影响
  3. 有时候一批 token 一起回来(特别是本地模型)

解决方案

方案一:前端缓冲(不推荐)

const buffer = [];
const flushInterval = 50;

setInterval(() => {
    if (buffer.length > 0) {
        const text = buffer.join('');
        appendToChat(text);
        buffer.length = 0;
    }
}, flushInterval);

这个方案的问题是会增加延迟,违背了流式的初衷。

方案二:平滑处理(推荐)

let lastAppendTime = 0;
const minInterval = 30;  // 最小间隔 30ms

ws.onmessage = (event) => {
    const data = JSON.parse(event.data);
    if (data.type === 'token') {
        const now = Date.now();
        const elapsed = now - lastAppendTime;

        if (elapsed >= minInterval) {
            appendToChat(data.content);
            lastAppendTime = now;
        } else {
            // 稍微延迟一下,避免太快
            setTimeout(() => {
                appendToChat(data.content);
                lastAppendTime = Date.now();
            }, minInterval - elapsed);
        }
    }
};

这样既能保持实时性,又不会忽快忽慢。

坑 3:高并发下内存飙升

上线后,并发量上来,服务端内存直接飙到 90%,差点 OOM。

排查过程

  1. memory_profiler 分析内存使用
  2. 发现 WebSocket 连接对象没有被释放
  3. 即使用户关闭了页面,服务端还保留着连接

解决方案

  1. 前端正确关闭 WebSocket:
window.addEventListener('beforeunload', () => {
    ws.close();
});
  1. 后端添加连接管理:
from collections import defaultdict
import weakref

# 存储活跃连接
active_connections = defaultdict(weakref.WeakSet)

@app.websocket("/ws/chat")
async def chat_websocket(websocket: WebSocket):
    await websocket.accept()
    active_connections['chat'].add(websocket)

    try:
        # ... 业务逻辑
    finally:
        active_connections['chat'].discard(websocket)
        await websocket.close()

# 定期清理无效连接
@app.on_event("startup")
async def setup_cleanup_task():
    async def cleanup():
        while True:
            await asyncio.sleep(300)  # 5 分钟清理一次
            for conn_set in active_connections.values():
                invalid = [ws for ws in conn_set if ws.client_state != 'connected']
                for ws in invalid:
                    conn_set.discard(ws)

    asyncio.create_task(cleanup())

坑 4:中文编码问题

传输中文时,有时候会出现乱码,或者 emoji 显示不出来。

原因

  1. JSON 默认不支持非 UTF-8 字符
  2. WebSocket 的数据格式处理不当

解决方案

# 后端:确保 JSON 编码正确
await websocket.send_json({
    "type": "token",
    "content": token
}, ensure_ascii=False)  # 关键:不转义非 ASCII 字符
// 前端:确保正确解码
ws.onmessage = (event) => {
    try {
        const data = JSON.parse(event.data);
        // ...
    } catch (e) {
        console.error('JSON 解析失败:', e);
    }
};

优化进阶

上面的实现已经能用了,但还有优化空间。

1. 首字延迟优化

TTFT 主要受这几个因素影响:

graph TD A[TTFT 总延迟] --> B[网络延迟] A --> C[队列等待] A --> D[模型推理] A --> E[序列化] B --> B1[往返时间] C --> C1[任务队列] D --> D1[模型加载] E --> E1[JSON 编码] style A fill:#f96 style D fill:#9f9,stroke:#333,stroke-width:2px

优化策略:

  1. 模型预加载:把模型常驻内存,避免每次加载
  2. 请求合并:同时来的多个请求可以合并处理
  3. 优先级队列:重要用户的请求优先处理
  4. 缓存:相同问题的答案可以缓存(注意时效性)

2. 批处理优化

如果用本地模型(比如 vLLM),可以启用批处理:

stream = await client.chat.completions.create(
    model="local-model",
    messages=[{"role": "user", "content": user_input}],
    stream=True,
    max_tokens=500,
    temperature=0.7,
    # vLLM 特定参数
    extra_body={
        "use_beam_search": False,
        "best_of": 1,
        "top_k": -1,
    }
)

3. 多模型并行

对于复杂任务,可以同时调用多个模型,选择最快返回的结果:

import asyncio

async def call_model(model_name, prompt):
    stream = await client.chat.completions.create(
        model=model_name,
        messages=[{"role": "user", "content": prompt}],
        stream=True
    )
    async for chunk in stream:
        if chunk.choices[0].delta.content:
            yield chunk.choices[0].delta.content

async def race_models(prompt):
    models = ["gpt-4", "gpt-3.5-turbo", "claude-3"]
    tasks = [call_model(model, prompt) for model in models]

    # 等待第一个模型返回第一个 token
    done, pending = await asyncio.wait(
        [asyncio.create_task(model.__anext__()) for model in tasks],
        return_when=asyncio.FIRST_COMPLETED
    )

    # 取消其他模型
    for task in pending:
        task.cancel()

    # 返回最快的结果
    if done:
        return done.pop().result()

4. 前端体验优化

除了技术上的优化,用户体验也很重要:

// 1. 骨架屏,首字返回前显示加载状态
function showSkeleton() {
    const chatContainer = document.getElementById('chat');
    chatContainer.innerHTML = '<div class="skeleton">正在思考...</div>';
}

// 2. 打字机光标效果
function addCursor() {
    const cursor = document.createElement('span');
    cursor.className = 'cursor';
    cursor.textContent = '|';
    document.getElementById('chat').appendChild(cursor);
}

// 3. 完成 1 秒后移除光标
ws.onmessage = (event) => {
    const data = JSON.parse(event.data);
    if (data.type === 'done') {
        setTimeout(() => {
            document.querySelector('.cursor')?.remove();
        }, 1000);
    }
};

结果对比

优化前后的对比:

指标优化前优化后提升
TTFT(首字延迟)2.3s420ms82%
TPOT(token 间隔)150ms60ms60%
总体延迟(用户感知)2.5s500ms80%
并发连接数100200020x
内存占用80%40%-50%

三项延迟指标放在一起看,首字延迟和总体感知延迟的降幅都在八成左右,流式改造的收益非常直观。

HTTP 改 WebSocket 流式后:TTFT、TPOT 与总体延迟优化前后对比(毫秒)

TPOT 从 150ms 压到 60ms,打字机效果终于稳定下来,用户反馈也从"卡死了"变成了"好快"。

总结

流式处理不是什么黑科技,核心就是"生成一点就返回一点",但真正要做好需要关注:

  1. 技术选型:WebSocket + 流式 API 是主流方案
  2. 架构设计:解耦前后端,中间用消息队列或 WebSocket 连接
  3. 踩坑准备:连接管理、内存泄漏、编码问题都要考虑
  4. 持续优化:TTFT、TPOT、并发能力都有优化空间

最后,流式处理虽然效果好,但也有适用场景:

  • ✅ 适合:实时对话、语音助手、代码补全
  • ❌ 不适合:批量生成、离线分析、需要完整结果的场景

如果你的场景需要实时性,流式处理值得投入。希望这篇文章能帮你少踩几个坑。

参考资料

版权声明: 本文首发于 指尖魔法屋-AI流式处理折腾手记https://blog.thinkmoon.cn/post/265-ai-streaming-latency-realtime-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!