详细讲解如何用 Python 实现生产级 OpenAI 兼容 API 流式客户端,涵盖 SSE、连接管理、超时重试等关键问题。
流式传输不仅仅是用户界面的技巧。它改变了服务如何管理连接、取消、重试和可观测性。等待完整响应的聊天 UI 在长时间生成期间可能显得已损坏;流式传输每个 token 的 UI 在客户端粗心重试时可能会留下半开连接和重复工作。
本指南展示了一个小型、面向生产的流式传输客户端,使用 Python 和 OpenAI 兼容端点。这些示例使用的是 2026-08-03 AIWave 公开定价目录中当前可见的模型标识符。费率仅作为已过期的参考:
在制定预算前查看 AIWave 定价页面。该端点是 OpenAI 兼容的,所以现有的 OpenAI SDK 集成可以保留其消息格式,同时使用 https://aiwave.live/v1 。
服务器发送事件(SSE)是一个长期存在的 HTTP 响应,其主体包含由空行分隔的事件。浏览器 EventSource API 是一个消费者,但 Python 服务可以用 HTTP 客户端消费相同的传输格式。一个事件通常看起来像 data: {json}\n\n;最终的 data: [DONE] 标记告诉客户端模型已完成。
SSE 是单向的。客户端无法通过同一流发送新消息。要取消生成,请关闭 HTTP 连接并在应用程序遥测中使取消可见。MDN SSE 参考描述了框架规则;你的提供商的 API 文档定义了每个事件内部的 JSON 字段。
OpenAI Python SDK 处理聊天完成的事件解析。在示例中保持密钥作为占位符,并通过部署密钥管理器加载真实密钥。
from openai import OpenAI
client = OpenAI(
base_url="https://aiwave.live/v1",
api_key="YOUR_API_KEY_HERE", # Create a key at https://aiwave.live/
timeout=60.0,
)
stream = client.chat.completions.create(
model="deepseek-v4-flash",
messages=[
{"role": "system", "content": "Answer in short, testable steps."},
{"role": "user", "content": "Explain why database indexes speed up reads."},
],
stream=True,
max_tokens=500,
)
for chunk in stream:
text = chunk.choices[0].delta.content or ""
print(text, end="", flush=True)
print()
第一个 token 是一个有用的延迟指标,但它与总请求延迟不同。同时记录 time_to_first_token 和 time_to_last_token。一个模型可以快速启动,但仍需要很长时间才能完成一个大的响应。
如果你的应用程序暴露自己的端点,不要在内存中缓冲提供商响应。在每个块到达时将其产出,并保留 SSE 内容类型。下面的示例使用 httpx 并向浏览器发送一个小的 JSON 信封。
import json
import httpx
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
app = FastAPI()
async def events(prompt: str):
payload = {
"model": "qwen3-coder-480b-a35b-instruct",
"messages": [{"role": "user", "content": prompt}],
"stream": True,
"max_tokens": 800,
}
headers = {"Authorization": "Bearer YOUR_API_KEY_HERE"}
async with httpx.AsyncClient(timeout=httpx.Timeout(60.0, connect=10.0)) as client:
async with client.stream(
"POST", "https://aiwave.live/v1/chat/completions", json=payload, headers=headers
) as response:
response.raise_for_status()
async for line in response.aiter_lines():
if not line.startswith("data: "):
continue
data = line[6:]
if data == "[DONE]":
yield "data: [DONE]\n\n"
return
event = json.loads(data)
delta = event.get("choices", [{}])[0].get("delta", {})
text = delta.get("content") or ""
if text:
yield "data: " + json.dumps({"text": text}) + "\n\n"
@app.get("/chat")
async def chat(prompt: str):
return StreamingResponse(events(prompt), media_type="text/event-stream")
在生产环境中,添加一个断开连接检查,以便浏览器离开时停止工作。FastAPI 的请求对象可以传入生成器中,并在块之间进行轮询。还要为路由关闭代理缓冲;否则 Nginx 或 CDN 可能会保留多个事件并破坏流式传输。
使用单独的连接和读超时。连接超时保护你免受上游死连接的影响;读超时保护你免受停止生成 token 的流的影响。正确的值取决于你的工作负载。交互式聊天可能使用 10 秒连接和 60 秒读取,而代码生成工作可能需要几分钟。
取消必须是幂等的。关闭客户端流是安全的,但立即重试可能会为一个用户操作创建两个生成。给每个请求一个内部 ID,并记录客户端是否断开连接、提供商是否正常结束,或你的超时是否终止了流。
重试在第一个 token 之前是最安全的。一旦输出到达用户,重试可能会重复文本或在模型调用工具时导致副作用两次。一个简单的策略是:
在第一个事件前重试连接失败和 429/5xx 响应。
使用带抖动的指数退避和较小的最大尝试次数。
永远不要在部分输出后盲目重试流;返回部分结果,并向调用方提供可重试状态。
对于工具调用代理,持久化工具调用 ID 并使下游操作幂等。流式传输改进了感知延迟,但它不能消除对分布式系统纪律的需要。
流式传输改变感知延迟,而不是 token 定价。在上述过期费率下,一个有 3000 个输入 token 和 600 个输出 token 的请求在 DeepSeek V4 Flash 上费用约为 $0.000865,在 Qwen3 Coder 上约为 $0.000576。公式是 (input_tokens × input_rate + output_tokens × output_rate) / 1,000,000;在财务报告前从 AIWave 定价刷新费率。
跟踪每个请求的这些字段:model、first-token 毫秒、total 毫秒、input tokens、output tokens、completion status、disconnect reason、retry count 和估计美元成本。除非已明确清理,否则将提示文本保持在日志外。一个请求 ID 足以将应用程序日志与提供商指标连接。
使用 AIWave 模型目录中的确切模型 ID。
设置单独的连接和读超时。
在流式传输代理路由上禁用缓冲。
客户端断开连接时停止上游工作。
仅在第一个 token 之前重试,使用有界退避。
将工具调用和其他副作用视为幂等。
在公开示例中保持 YOUR_API_KEY_HERE,并通过密钥管理器轮转真实密钥。
聊天完成文档涵盖请求参数和流式传输行为。从一个小的评估集开始,分别测量 first-token 和完成延迟,然后使用观察到的生产分布而不是单个基准来调整超时。