AI API 流式输出(SSE)怎么正确实现:从 SDK 判空到 nginx 缓冲,一次讲透

SSE 到底是什么:一条 HTTP 长连接,不是 WebSocket

流式输出的底层协议叫 SSE(Server-Sent Events,服务器推送事件)。你发一个普通 POST 请求,服务器不把响应立刻收尾,而是把连接挂着,模型生成几个字就往回吐几个字,直到说完才关闭。看响应头 Content-Type: text/event-stream 就能判断:回来的是一串碎片,不是完整 JSON,解析思路得整段换掉。

数据帧格式:data: 开头,空行才是分隔

SSE 每个数据帧都是纯文本,格式固定成两行:一行 data: {...},加一个空行。空行才是消息分隔符,不是回车也不是分号。少了它,EventSource 和多数解析库都不认为这一帧结束,数据到了界面也不渲染。

带上 stream: true 之后,上游返回的字节大概长这样:

解析器的正确读法是"攒到一个空行就交出一帧",不是"读到一行就当成一条消息"。

[DONE] 是结束标记,它不是 JSON

[DONE] 是 OpenAI 兼容协议约定的流结束信号。它不是合法 JSONjson.loads("[DONE]") 会抛 json.JSONDecodeError: Expecting value: line 1 column 1 (char 0)。有人习惯先把每行都 json.loads 一遍再判断,结果每次正常结束都报错。顺序应该是先判 [DONE]break,再解析。

为什么不用 WebSocket

WebSocket 要先握手升级,双方还得维护连接状态、心跳、重连、背压,复杂度是 SSE 的好几倍。聊天本来就是单向的——我发一次提问,模型回一段文本,用不上双向通道。SSE 就是普通 HTTP 响应,代理、网关、浏览器都支持。Chatbox、Cherry Studio、NextChat、Cline、RooCode、KiloCode、LobeChat 内部用的都是同一套 SSE 解析逻辑。

为什么值得为流式多写这些代码

答案就一个字:等。思考型模型(名字带 thinking 的那类)在吐出开头的字之前要先跑完一段推理链,首字延迟普遍落在 8 到 20 秒;非流式更难受,得等整段话生成完才返回。一个 800 字的回答,非流式下要盯着空白页面等 25 秒以上,很多人以为程序死了,刷新重来——钱花两次,字一个没看到。

非流式下你等的是全文,不是首字

非流式的时间线是:发送请求 → 模型推理 → 生成全文 → 一次返回,前端测到的延迟等于总耗时。流式把这条线拆开了:首字仍要等 8 到 20 秒,之后每 30 到 80 毫秒来一批字。总耗时差不多,感受却是"它在想"和"它死了"的区别。

按次计费下,流式不额外收钱

有人以为流式更贵,理由是"长连接占服务器资源",其实不是。小鱼API 以按次计费为主:一次请求一个固定价,输入输出长度都不影响扣费,stream: truestream: false 扣的是同一个价。deepseek-v3.2-thinking 0.049 元一次,流式也是 0.049 元。另有按量计费与无限卡套餐可选,但流式这个开关本身不多收钱。

流式与非流式怎么选

维度流式(stream)非流式
首字延迟8-20 秒(思考型模型)等于全文生成耗时
总耗时与非流式接近,略高一点略低
内存占用常数级,边收边处理需缓存完整响应
实现复杂度高:分帧、判空、残余缓冲低:一次 json.loads
超时风险低:有数据流动就不算卡死高:长文本易被网关掐断
适用场景对话界面、批量长文本结构化抽取、短分类

openai SDK 流式:判空是保命操作

openai 这个 Python SDK 写流式最省事,坑全在边界上。下面这段可以直接跑,Key 放进环境变量 XYUAI_API_KEY。主入口 https://xyuapi.top/v1,备用 https://xyuai.cc/v1,协议一致;Gemini 原生协议另有 /v1beta 端点。

完整代码:能直接跑的流式调用

import os
from openai import OpenAI, APIStatusError, APITimeoutError

client = OpenAI(
    api_key=os.environ["XYUAI_API_KEY"],
    base_url="https://xyuapi.top/v1",   # 备用:https://xyuai.cc/v1
    timeout=300.0,                      # 流式下这是"块间隔"上限,不是总时长
    max_retries=2,
)

def stream_chat(prompt: str, model: str = "deepseek-v3.2-thinking") -> str:
    """流式调用,返回拼接后的完整文本。空值、用量帧、异常都在这里处理干净。"""
    pieces = []
    resp = None
    try:
        resp = client.chat.completions.create(
            model=model,
            messages=[{"role": "user", "content": prompt}],
            stream=True,
            stream_options={"include_usage": True},   # 末尾多回一帧用量
        )
        for chunk in resp:
            # 边界一:用量帧的 choices 是空列表,取 [0] 会 IndexError
            if not getattr(chunk, "choices", None):
                usage = getattr(chunk, "usage", None)
                if usage:
                    print(f"\n[用量] total={usage.total_tokens}")
                continue
            delta = chunk.choices[0].delta
            # 边界二:开头那帧只有 role,content 是 None
            text = getattr(delta, "content", None)
            if not text:
                continue
            pieces.append(text)
            print(text, end="", flush=True)   # flush 不能省,否则看起来像卡住
    except APIStatusError as e:
        # 429 是限速,5xx 是上游抖动,都按可重试处理
        detail = e.response.text[:200] if e.response is not None else ""
        print(f"\n[失败] status={e.status_code} body={detail}")
        raise
    except APITimeoutError as e:
        print(f"\n[超时] 两个数据块之间超过 300 秒没来数据: {e}")
        raise
    except ConnectionError as e:
        print(f"\n[断连] {type(e).__name__}: {e}")
        raise
    finally:
        # 成功、报错、被中断都要关,避免连接池里挂着半截流
        if resp is not None:
            resp.close()
    return "".join(pieces)

if __name__ == "__main__":
    total = stream_chat("用三句话说明 SSE 的帧格式")
    print(f"\n--- 共 {len(total)} 字 ---")

delta.content 可能是 None,不判空就崩

三个必须判空的位置,漏一个就等着半夜看报错:

stream_options 怎么拿 token 用量

stream_options={"include_usage": True} 会让上游在结束前补推一帧,里面是 prompt_tokenscompletion_tokenstotal_tokens。不开这个参数,流式响应里根本没有 usage 字段,成本统计无从下手。这一帧的 choices 是空数组,所以必须先判 choices 再取用量。

requests 手工解析:裁前缀、遇 DONE 退出

不用 SDK 时,常见写法是 requests.post(..., stream=True)for chunk in r.iter_content() 然后 chunk.decode()。这条路上有两个必踩的坑,下面这段代码一起绕过去。

完整代码:手工解析 SSE

import json
import os
import time
import requests

URL = "https://xyuapi.top/v1/chat/completions"   # 备用:https://xyuai.cc/v1
HEADERS = {
    "Authorization": f"Bearer {os.environ['XYUAI_API_KEY']}",
    "Content-Type": "application/json",
    "Accept": "text/event-stream",
}

def sse_chat(prompt: str, model: str = "deepseek-v3.2-thinking", read_timeout: float = 300.0) -> str:
    payload = {
        "model": model,
        "messages": [{"role": "user", "content": prompt}],
        "stream": True,
    }
    pieces, t0, first_at = [], time.time(), None
    # timeout 传元组 (连接超时, 读取超时);读取超时管的是两个数据块之间的间隔
    with requests.post(URL, headers=HEADERS, json=payload,
                       stream=True, timeout=(10.0, read_timeout)) as r:
        r.raise_for_status()
        ctype = r.headers.get("Content-Type", "")
        if "text/event-stream" not in ctype:
            raise RuntimeError(f"上游没按 SSE 返回,Content-Type={ctype!r}")

        # iter_lines(decode_unicode=True) 会自己攒字节直到换行,字符边界天然对齐
        for raw in r.iter_lines(decode_unicode=True):
            if raw is None:            # 空行 = 一帧结束,没内容就跳过
                continue
            line = raw.strip()
            if not line or line.startswith(":"):
                continue               # 空行、以及以冒号开头的保活注释
            if not line.startswith("data:"):
                continue               # event: / id: / retry: 这些字段用不上
            data = line[5:].strip()    # 关键:裁掉 "data:" 前缀
            if data == "[DONE]":
                break                  # 结束标记,先判它再解析
            try:
                obj = json.loads(data)
            except json.JSONDecodeError:
                continue               # 半条消息或心跳碎片,丢掉而不是崩
            choices = obj.get("choices") or []
            if not choices:
                continue
            text = (choices[0].get("delta") or {}).get("content")
            if not text:
                continue
            if first_at is None:
                first_at = time.time() - t0
            pieces.append(text)
            print(text, end="", flush=True)

    print(f"\n[首字 {first_at:.2f}s | 总计 {time.time() - t0:.2f}s]")
    return "".join(pieces)

if __name__ == "__main__":
    try:
        sse_chat("解释一下 SSE 里空行的作用")
    except requests.exceptions.ReadTimeout:
        print("读取超时:上游卡住了,可以拿已收到的部分结果先落盘")
    except requests.exceptions.HTTPError as e:
        print(f"HTTP 错误: {e.response.status_code} {e.response.text[:200]}")

为什么不能对每行裸调 json.loads

流式响应里下面这几种行都是正常的,裸调就会抛 json.JSONDecodeError,把整个循环打断:空行本身、以冒号开头的保活注释、data: [DONE] 结束标记、上游偶发截断的半条消息。处理原则只有一条:解析失败就 continue,不要 break。流式场景里"丢掉一帧"可以接受,"整段崩掉"不行。

中文乱码:别按字节块 decode

一个汉字在 UTF-8 里占 3 个字节,而 TCP 分片和 iter_content 的 chunk 边界跟字符边界毫无关系,完全可能把"流"字的 3 个字节拆成 2+1。你在循环里对每个 chunk 单独 .decode("utf-8"),就会撞上这样的报错:

UnicodeDecodeError: 'utf-8' codec can't decode byte 0xe6 in position 0: unexpected end of data

有人改成 .decode("utf-8", errors="ignore"),报错没了,但那截不完整的字节被静默丢弃,回答里莫名少字,而且只在偶发分片时出现,上线后靠用户截图才发现,更难查。正确做法是按行切:iter_lines(decode_unicode=True) 自己维护缓冲,遇到换行符才交出一整行。想自己实现思路一样——攒一个 bytes 缓冲,只在找到 \n 的位置切分。

FastAPI 转发给浏览器:响应头比函数更容易错

前端要调自己的后端(而不是把 Key 暴露在浏览器里)时,后端得当一次二传手:把上游的 SSE 边收边转给浏览器。函数本身三行写完,熬夜的往往是响应头。

完整代码:边收边转

import json
import os
import httpx
from fastapi import FastAPI
from fastapi.responses import StreamingResponse

app = FastAPI()
UPSTREAM = "https://xyuapi.top/v1/chat/completions"
API_KEY = os.environ["XYUAI_API_KEY"]

SSE_HEADERS = {
    "Cache-Control": "no-cache",
    "Connection": "keep-alive",
    "X-Accel-Buffering": "no",   # 告诉 nginx:这条响应别做代理缓冲
}

@app.post("/api/chat/stream")
async def chat_stream(payload: dict):
    payload["stream"] = True

    async def gen():
        # connect 10 秒够;read 给 300 秒,它管的是两个数据块之间的最大间隔
        timeout = httpx.Timeout(connect=10.0, read=300.0, write=30.0, pool=10.0)
        async with httpx.AsyncClient(timeout=timeout) as client:
            try:
                async with client.stream(
                    "POST", UPSTREAM, json=payload,
                    headers={"Authorization": f"Bearer {API_KEY}",
                             "Content-Type": "application/json",
                             "Accept": "text/event-stream"},
                ) as up:
                    if up.status_code != 200:
                        body = (await up.aread()).decode("utf-8", "replace")[:300]
                        err = {"error": f"上游 {up.status_code}: {body}"}
                        yield f"data: {json.dumps(err, ensure_ascii=False)}\n\n"
                        return
                    async for line in up.aiter_lines():
                        if not line.strip():
                            continue
                        # 补回空行:少一个空行,浏览器就不认这一帧
                        yield line + "\n\n"
                yield "data: [DONE]\n\n"      # 放在 try 里,不能放 finally
            except httpx.ReadTimeout:
                yield 'data: {"error": "上游读超时,两个数据块之间超过 300 秒"}\n\n'
            except httpx.RemoteProtocolError:
                yield 'data: {"error": "上游连接中途断开"}\n\n'

    return StreamingResponse(gen(), media_type="text/event-stream", headers=SSE_HEADERS)

三个响应头各管什么

media_type 必须是 text/event-stream;写成 application/jsonEventSource 会直接拒绝,fetch 也拿不到流式分帧。

nginx 反代缓冲:回答一次性蹦出来的元凶

现象好认:本地 uvicorn 一切正常,挂到 nginx 后面就变成转圈 20 秒、整段回答一次性蹦出来——nginx 默认开代理缓冲,把上游碎块攒成整包再转发。在对应 location 里补一行 proxy_buffering off;,再把 proxy_read_timeout 从默认 60 秒调到 300 秒(长回答会被它掐断,前端收到 504)。X-Accel-Buffering: no 和它不冲突:前者让后端替 nginx 做决定,后者把规则写死在网关侧,两个都加更稳。

客户端断开要止损

用户在生成过程中关掉页面,浏览器会断开那条 fetch 连接。后端不处理的话,它会继续从上游拉数据、继续计费,而结果永远没人看。做法是把生成逻辑放进 try/finally,用 async with httpx.AsyncClient() 管理上游连接——生成器被回收或连接断开时,Python 会往里抛 GeneratorExit,清理代码一定执行。别在 finallyyield,那时生成器已经进入关闭状态,会报 RuntimeError: async generator ignored GeneratorExitdata: [DONE] 要写在 try 里、循环结束之后。

浏览器端解析:一个 chunk 可能只有半条消息

浏览器拿到的 ReadableStream 分片完全看网络心情:一次 reader.read() 可能给你三条完整消息,也可能只给半条,{"choices":[{"delta":{"cont 到这儿就断了。所以必须有残余缓冲区。

完整代码:fetch 加 ReadableStream

async function streamChat(prompt, onDelta, { signal } = {}) {
  const res = await fetch('/api/chat/stream', {
    method: 'POST',
    headers: { 'Content-Type': 'application/json' },
    body: JSON.stringify({
      model: 'deepseek-v3.2-thinking',
      messages: [{ role: 'user', content: prompt }],
      stream: true,
    }),
    signal,
  });
  if (!res.ok) throw new Error(`HTTP ${res.status}`);
  if (!res.body) throw new Error('响应没有 body,可能被代理缓冲了');

  const reader = res.body.getReader();
  const decoder = new TextDecoder('utf-8');
  let buffer = '';                      // 残余的半条消息留在这里
  let full = '';
  try {
    while (true) {
      const { done, value } = await reader.read();
      if (done) break;
      // stream: true 让解码器自己留住不完整的字节序列
      buffer += decoder.decode(value, { stream: true });

      // SSE 以空行分帧,所以按 \n\n 切,不能按单个 \n 切
      const events = buffer.split('\n\n');
      buffer = events.pop();            // 末尾那截可能不完整,留到下一轮

      for (const evt of events) {
        for (const line of evt.split('\n')) {
          if (!line.startsWith('data:')) continue;
          const data = line.slice(5).trim();
          if (data === '[DONE]') return full;
          try {
            const obj = JSON.parse(data);
            const text = obj?.choices?.[0]?.delta?.content;
            if (text) { full += text; onDelta(text); }
          } catch (e) {
            continue;                   // 半条 JSON 或心跳,跳过就好
          }
        }
      }
    }
  } finally {
    reader.cancel().catch(() => {});     // 用户中断、网络异常都要放掉连接
  }
  return full;
}

TextDecoder({stream: true}) 为什么不能省

这是后端乱码坑的浏览器版。new TextDecoder().decode(chunk) 默认把每次调用当成相互独立的完整字节序列,一旦碰上被切开的汉字,就吐出 (U+FFFD 替换字符)。传 {stream: true} 之后,解码器会把不完整的字节留在内部,等下一次 decode 补齐再一起输出。顺序也要留意:先解码出字符串,再拼 buffer、再切帧,反过来先切字节就会切在字符中间。

分帧要按 \n\n,不是 \n

空行才是帧边界。用 \n 切,data: 行会和下一帧粘在一起;用 \n\n 切,末尾那段可能不完整,所以要 pop() 出来留到下一轮。代码里的 buffer = events.pop() 干的就是这件事,漏掉它,长回答会随机丢掉尾部几个字,短回答又完全看不出来。

超时怎么设:读超时指的是块间隔

流式下超时的含义和非流式完全不同,配错了要么频繁误杀正常请求,要么请求挂了半小时也不报警。

连接 10 秒,读取 180 到 300 秒

长回答配保活心跳

思考型模型在推理阶段可能 30 秒不出任何文本 token,这段时间连接上一个数据帧都没有,读取超时就在倒计时。两个办法:把读取超时放宽到 300 秒;在自己的转发层每 15 秒补一个注释帧(一行以冒号开头的空内容),让连接上始终有字节流动。

排错表与成本账

现象 → 原因 → 解决

现象原因解决
回答一次性蹦出来nginx 代理缓冲把碎块攒成整包proxy_buffering off;X-Accel-Buffering: no
卡住不动,20 秒后整段出现中间层没认 text/event-stream检查 media_type 与上游 Content-Type
IndexError: list index out of range用量帧的 choices 是空列表,却直接取了 [0]choices[0] 前判空
TypeError: expected str instance开头那帧 delta.contentNone,被 append 进去取到 content 后判空再加
中文变问号或 \ufffd按字节块 decode,汉字的 3 个字节被切断按行切;浏览器端用 TextDecoder({stream:true})
少字但没有任何报错用了 errors="ignore",坏字节被静默丢弃去掉 ignore,改成攒 buffer 按行切
连接中途断开网关 proxy_read_timeout 默认 60 秒到期调到 300 秒,并加 15 秒保活注释帧
没收到 [DONE]上游异常退出,或循环没走到结束就 return先判 [DONE],异常分支也补发结束帧
关掉页面后额度还在掉客户端断开了,后端仍在拉上游try/finally 里关掉上游连接

500 次批量任务的账

deepseek-v3.2-thinking 的 0.049 元每次算,500 次生成任务跑一遍是 500 × 0.049 = 24.5 元。非流式下,长回答容易撞上网关或客户端超时被掐断,被中断的请求拿不到可用结果,只能整条重跑。假设中断率 20%,多花的是 100 × 0.049 = 4.9 元,账单从 24.5 元变成 29.4 元,成本上浮 20%。流式的价值就在这儿:文字陆续到达,即使中途断流,前面 80% 的内容也已落盘,补跑只要续上尾部,不用把整条问题从头问一遍。批量任务的墙钟时间也一起降:500 条并发跑,边生成边写文件,不用等每条都完整返回。

几种常用模型的单价摆在一起,方便按预算挑:

模型单价(元/次)500 次合计
gemini-2.5-pro0.03115.5 元
deepseek-v3.2-thinking0.04924.5 元
deepseek-v4-flash-thinking0.0525 元
kimi-k2.60.0945 元
claude-sonnet-4-5-thinking0.0945 元
gpt-5.50.2100 元
claude-opus-4-6-thinking0.25125 元

小鱼API 的最低充值 7 元,支付宝或微信直接付,不需要海外信用卡。按次计费下上面这些数字就是最终价,提示词写多长都不影响单次价格。

上线前的自检清单

相关阅读

🚀 想要立即使用?来小鱼API体验全系列AI模型

查看全部产品

支持Gemini / Claude / GPT / Grok / DeepSeek · 国内直连 · 支付宝/微信