AI 工程基础体系 · 第 15/100 篇。内容覆盖机器学习、深度学习与生成式 AI;模型、数据、评测、权限和成本会作为同一生产系统处理。
LLM API 工程:客户端、流式输出、取消、重试、限流和兼容层
LLM API 工程不是“调用一个模型接口并读取字符串”。在生产系统中,一次生成请求同时受到模型能力、上下文长度、输入输出 Token、网络连接、并发配额、用户取消、供应商协议差异和成本预算的影响。
一个可用的 LLM 客户端至少要回答以下问题:
- 请求如何建立、复用和关闭?
- 非流式响应与流式响应的生命周期有什么不同?
- 用户断开连接时,如何停止上游生成,避免继续消耗 Token?
- 哪些错误可以重试,哪些错误重试只会放大故障?
- 如何同时控制请求数、并发数和 Token 吞吐量?
- 如何在不同模型供应商、不同 SDK 和 OpenAI-compatible 服务之间保持稳定的应用接口?
- 如何把延迟、Token、取消、重试和成本记录到同一条 Trace 中?
本文中的“LLM API”指通过 HTTP 或 SDK 调用生成式模型的远程推理接口。模型本身可能是闭源托管模型、云厂商模型,也可能是基于 Transformers、vLLM、TGI 或其他推理服务器部署的开源模型。API 工程关注的是调用边界,不假设所有模型在能力、协议和错误语义上相同。
一、先建立正确的调用模型
1. 客户端不是一次性函数,而是资源管理器
一个 LLM 请求通常经历以下阶段:
- 构造输入消息、工具定义、结构化输出约束和模型参数。
- 从客户端连接池获取或建立 TCP/TLS 连接。
- 发送 HTTP 请求。
- 等待服务端接受请求并开始推理。
- 接收完整响应,或持续接收流式事件。
- 解析响应、记录 Token 和耗时。
- 释放响应体、连接和本地任务资源。
因此,客户端对象通常应在进程或应用生命周期内复用,而不是每个请求都重新创建。复用客户端可以复用连接池、TLS 会话和 DNS 结果,减少握手开销;但客户端通常不是无条件线程安全的,具体并发保证要以 SDK 文档为准。异步应用应使用异步客户端,不能在事件循环中直接调用阻塞式同步客户端。
客户端生命周期可以抽象为:
创建配置
└── 创建 HTTP/SDK 客户端
├── 多次发送请求
├── 复用连接池
└── 应用退出时关闭客户端
下面是一个基于 Python 异步 SDK 的最小非流式示例。示例采用 OpenAI Python SDK 的 Responses API 形式;SDK 版本、模型名和具体参数可能随供应商变化,实际使用时应以对应版本的官方文档为准。
import asyncio
import os
from openai import AsyncOpenAI
async def main() -> None:
client = AsyncOpenAI(
api_key=os.environ["OPENAI_API_KEY"],
timeout=30.0,
max_retries=0, # 将重试交给应用层,避免 SDK 与应用层重复重试
)
try:
response = await client.responses.create(
model=os.environ["OPENAI_MODEL"],
input=[
{
"role": "user",
"content": "用两句话解释什么是幂等性。",
}
],
)
print(response.output_text)
finally:
await client.close()
if __name__ == "__main__":
asyncio.run(main())
这里有三个重要边界:
timeout=30.0通常表示 HTTP 请求等待上限,不一定等于“模型必须在 30 秒内生成完”。max_retries=0不是说系统永不重试,而是避免 SDK 内部重试与应用层重试叠加。response.output_text是 SDK 对响应结构的便利封装;底层响应可能包含多个输出项、工具调用、拒答、元数据和 Token 用量,不能假设所有响应都只有一段文本。
2. 总超时、连接超时、读取超时要分开
“超时”不是单一事件。至少应区分:
- 连接超时:无法在规定时间内建立连接。
- 写入超时:请求体无法发送完,例如网络阻塞。
- 首字节或首事件超时:服务端迟迟没有开始返回。
- 读取空闲超时:流式连接已经建立,但连续一段时间没有新事件。
- 总耗时超时:整个请求超过业务允许的最大时间。
- 客户端取消:用户或上游服务主动不再需要结果。
流式输出尤其不能只设置一个总 HTTP 超时。模型可能在生成较长答案时,连续事件之间暂时存在间隔;如果空闲超时过短,会把正常生成误判为故障。反过来,如果没有空闲超时,半断开的连接可能长期占用并发槽位。
工程上可以使用两个时间约束:
其中 是请求总预算, 是相邻流事件之间的时间间隔, 是允许的最大空闲时间。两者分别防止“请求整体过慢”和“连接失活但不退出”。
二、流式输出:传输层、事件层和文本层不是一回事
1. 流式输出的本质是增量事件
非流式请求通常在服务端生成完成后返回一个完整 JSON。流式请求则让服务端在生成过程中持续发送事件。常见实现是 SSE(Server-Sent Events),其传输格式大致如下:
event: response.output_text.delta
data: {"delta":"你好"}
event: response.output_text.delta
data: {"delta":",世界"}
event: response.completed
data: {"id":"resp_123"}
SSE 是 HTTP 之上的文本事件协议。它不等于“每个事件就是一个 Token”,原因包括:
- 服务端可能按内部批次发送多个 Token。
- 一个 Token 可能对应一个 UTF-8 字节序列的一部分或多个字符。
- 事件可能承载文本增量、工具调用参数、状态变化、错误和完成信息。
- 代理、网关或 SDK 可能重新缓冲事件。
因此,应用应消费“事件”和“文本增量”,而不是自行假设事件边界就是 Token 边界。
2. 流式请求的状态机
一个可靠的流式请求至少包含以下状态:
stateDiagram-v2
[*] --> Created
Created --> Connecting
Connecting --> Streaming
Connecting --> Failed
Streaming --> Streaming: 接收增量事件
Streaming --> Completed: 收到完成事件
Streaming --> Failed: 上游错误/协议错误
Streaming --> Cancelled: 客户端取消
Streaming --> TimedOut: 总超时或空闲超时
Completed --> [*]
Failed --> [*]
Cancelled --> [*]
TimedOut --> [*]
Streaming 状态不能仅由“已经收到文本”定义。某些模型调用会先产生工具调用事件,再产生文本;有些请求可能只产生结构化输出或拒答。因此状态解析应依据事件类型和终止事件,而不是把所有非文本事件丢掉后等待字符串结束。
3. 一个可运行的异步流式示例
import asyncio
import os
from openai import AsyncOpenAI
async def stream_answer() -> None:
client = AsyncOpenAI(
api_key=os.environ["OPENAI_API_KEY"],
timeout=60.0,
max_retries=0,
)
try:
stream = await client.responses.create(
model=os.environ["OPENAI_MODEL"],
input="流式输出和一次性返回有什么区别?",
stream=True,
)
async for event in stream:
event_type = getattr(event, "type", "")
if event_type == "response.output_text.delta":
# delta 是本次事件新增的文本,而不是完整答案
print(event.delta, end="", flush=True)
elif event_type == "response.completed":
print("\n[completed]")
elif event_type == "response.failed":
print("\n[failed]", event)
break
finally:
await client.close()
asyncio.run(stream_answer())
运行前需要安装 SDK,并设置 OPENAI_API_KEY 与 OPENAI_MODEL。预期输出类似:
流式输出和一次性返回有什么区别?
[completed]
实际输出会被拆成多个增量事件。print(..., flush=True) 的作用是及时刷新本地标准输出;它不能改变服务端的发送粒度,也不能保证浏览器端立即显示。如果中间还有 Nginx、网关或 WebSocket/SSE 转发层,还必须检查这些组件是否关闭了响应缓冲。
4. 应用层转发时必须处理背压
“背压”是指下游消费速度低于上游生产速度。对于 LLM 流式响应,典型路径是:
模型服务 → API 客户端 → 应用服务器 → 浏览器/移动端
如果浏览器网络很慢,而应用仍不断从上游读取并把数据放进无限队列,内存会增长;如果应用完全不读取上游,连接可能积压,最终触发超时。
一个合理的转发实现应使读取和发送受到同一个生命周期控制:
async for event in upstream:
chunk = encode_for_browser(event)
await downstream.send(chunk)
这里的 await downstream.send 形成了自然背压:下游发送受阻时,应用不会无限快地读取上游。若必须解耦生产者和消费者,应使用有界队列:
queue = asyncio.Queue(maxsize=64)
有界队列满时,生产者必须等待或主动取消上游,而不是继续分配内存。
5. 文本拼接、结构化输出和工具调用不能混为一谈
对于纯文本,常见做法是:
full_text = ""
async for event in stream:
if event.type == "response.output_text.delta":
full_text += event.delta
但生产代码通常应使用列表收集增量,再在结束时拼接:
parts = []
async for event in stream:
if event.type == "response.output_text.delta":
parts.append(event.delta)
full_text = "".join(parts)
反复对不可变字符串使用 += 可能造成多次复制,长输出时效率更差。
如果输出是 JSON,不能在每个增量到达时都假设它是完整 JSON。下面的片段在流式阶段通常是非法 JSON:
{"name":"Alice","age":
应在完整输出结束后解析,或使用供应商提供的结构化事件和校验机制。若模型输出工具调用参数,还应区分:
- 工具名增量;
- 参数 JSON 增量;
- 工具调用完成;
- 工具执行结果;
- 第二轮模型生成。
把工具参数的半截 JSON 直接交给业务函数,会导致解析错误或危险的部分执行。
三、取消:停止客户端等待,不等于停止服务端推理
1. 取消有多个层级
“取消请求”至少可能指:
- 业务层取消:用户点击停止、页面关闭或上游任务失效。
- 应用任务取消:Python 的
asyncio.Task被取消。 - HTTP 连接取消:客户端关闭响应流或断开连接。
- 供应商任务取消:通过供应商提供的取消接口通知服务端停止生成。
- 模型计算取消:推理服务器真正中止 GPU 上的计算。
前两个层级由应用控制,第三个层级依赖 HTTP 客户端,第四个层级依赖 API 能力,第五个层级还依赖供应商内部实现。客户端关闭连接并不能规范性地保证服务端已经停止计算。
2. 用户断开时的取消路径
典型路径如下:
sequenceDiagram
participant U as 用户浏览器
participant A as 应用服务器
participant C as LLM 客户端
participant M as 模型服务
U->>A: 建立流式请求
A->>C: 创建异步生成任务
C->>M: HTTP 流式请求
M-->>C: 增量事件
C-->>A: 转发增量
A-->>U: 显示文本
U-xA: 断开连接
A->>C: cancel task / close stream
C-xM: 关闭 HTTP 连接
M-->>C: 可能继续、可能停止、可能异步取消
x 表示连接被关闭,而不是服务端一定完成了计算中止。因此取消后仍可能产生少量计费 Token,具体取决于供应商计费规则和服务端取消时点。
3. 使用 asyncio 传播取消
下面示例展示应用层的基本取消逻辑:
import asyncio
import os
from openai import AsyncOpenAI
async def consume_stream(client: AsyncOpenAI) -> str:
parts = []
stream = await client.responses.create(
model=os.environ["OPENAI_MODEL"],
input="写一段较长的说明。",
stream=True,
)
async for event in stream:
if event.type == "response.output_text.delta":
parts.append(event.delta)
print(event.delta, end="", flush=True)
return "".join(parts)
async def main() -> None:
client = AsyncOpenAI(
api_key=os.environ["OPENAI_API_KEY"],
timeout=120.0,
max_retries=0,
)
task = asyncio.create_task(consume_stream(client))
try:
await asyncio.sleep(2)
# 模拟用户点击“停止”
task.cancel()
try:
await task
except asyncio.CancelledError:
print("\n[client cancelled]")
finally:
await client.close()
asyncio.run(main())
关键点是必须 await task,这样取消异常才会真正传播并完成任务清理。只调用 task.cancel() 而不等待任务,可能留下未完成任务、未关闭响应体或异常未处理。
如果 SDK 暴露显式关闭流的方法,应在 finally 中调用。若 SDK 没有公开关闭方法,则应依赖任务取消和 HTTP 客户端上下文管理;不能臆造一个不存在的 close_stream() API。
4. 取消的幂等性
取消操作通常应设计为幂等:
用户可能重复点击停止,网络重试也可能重复发送取消请求。服务端收到第一次取消后,第二次应返回“已完成”“已取消”或“资源不存在”等可接受状态,而不应破坏其他请求。
应用层还应记录取消原因:
user_cancelledclient_disconnecteddeadline_exceededupstream_failedbudget_exceeded
这些原因对成本分析和用户体验诊断不同。把所有中止都记成“模型错误”会掩盖真实问题。
四、重试:先判断故障是否可重试,再决定等待多久
1. 错误分类比重试代码更重要
常见错误可按可重试性分组:
| 错误类型 | 典型原因 | 通常是否重试 |
|---|---|---|
| DNS、连接失败 | 网络或服务不可达 | 可以有限重试 |
| 连接/读取超时 | 网络抖动或服务过慢 | 视请求是否已到达服务端 |
| HTTP 408 | 请求超时 | 通常可重试 |
| HTTP 429 | 限流或配额不足 | 按服务端提示延迟 |
| HTTP 500、502、503、504 | 服务端暂时故障 | 可以有限重试 |
| HTTP 400 | 参数、格式或上下文错误 | 不应原样重试 |
| HTTP 401、403 | 认证或权限错误 | 不应重试 |
| 内容安全拒绝 | 输入或输出违反策略 | 通常不应原样重试 |
| 业务校验失败 | 模型输出不符合格式 | 应修复提示或降级,不是 HTTP 重试 |
HTTP 状态码只是第一层信号。响应体中的错误类型、请求 ID 和供应商特定错误码往往更准确。
2. 流式请求的重试更危险
非流式请求在客户端没有收到响应时,无法确定服务端是否已经执行了请求。例如:
客户端发送成功
服务端完成模型生成
服务端准备返回
连接在返回前断开
客户端看到超时
此时客户端重试,可能导致同一个请求执行两次。对于纯文本生成,重复通常只是浪费成本;对于模型调用工具、发送邮件、创建订单等场景,重复可能造成真实副作用。
所以需要区分:
- 生成操作:通常无外部副作用,但可能重复计费。
- 工具执行:可能产生外部副作用,必须使用幂等键或业务去重。
- 流式输出:已经向用户展示了部分结果后再重试,不能简单拼接两次答案。
如果流式连接在输出中途断开,应用有三种策略:
- 标记失败,要求用户重新生成;
- 使用相同请求标识重新请求,并覆盖旧输出;
- 继续请求并尝试恢复,但需要服务端支持可恢复游标或响应 ID。
第三种不能通过“把已输出文本放入新 Prompt”简单实现,因为模型可能重新生成不同内容,也会增加输入 Token。
3. 指数退避和抖动
常见退避时间为:
其中:
- 是第 次重试;
- 是初始等待时间;
- 是最大等待时间;
- 是范围为 到 的随机抖动。
随机抖动可以避免大量客户端在同一时刻同时重试,形成“惊群”。
若响应包含 Retry-After,通常应优先遵守服务端给出的等待时间,但仍应受本地最大等待和总截止时间约束。
import asyncio
import random
async def retry_delay(attempt: int, base: float = 0.5,
cap: float = 8.0, jitter: float = 0.2) -> float:
exponential = min(cap, base * (2 ** attempt))
return exponential + random.uniform(0, jitter)
async def retry_loop(operation, max_attempts: int = 4):
for attempt in range(max_attempts):
try:
return await operation()
except Exception as exc:
retryable = is_retryable(exc) # 需按 SDK 异常和 HTTP 状态实现
if not retryable or attempt == max_attempts - 1:
raise
delay = await retry_delay(attempt)
await asyncio.sleep(delay)
def is_retryable(exc: Exception) -> bool:
# 这里只是结构示意;生产代码应检查具体 SDK 异常类型和状态码
return True
这个示例故意没有把所有异常都判为可重试。真正实现中,is_retryable 必须明确排除认证失败、参数错误、上下文过长和策略拒绝。
4. 重试预算必须独立于请求预算
如果每个原始请求最多重试 3 次,流量高峰时实际请求数可能接近原来的 4 倍。重试应受以下约束:
- 单请求最大尝试次数;
- 单请求总时间预算;
- 服务实例级重试并发上限;
- 全局错误率熔断;
- 供应商配额和成本预算。
设原始请求速率为 ,平均尝试次数为 ,则上游实际请求速率近似为:
当 429 或 503 增多时, 可能上升,进一步推高 ,形成正反馈。限流和重试必须一起设计,否则“重试恢复故障”会变成“重试制造故障”。
五、限流:请求数、并发数和 Token 数是三个不同资源
1. RPM、TPM 和并发限制
LLM 服务常见三类限制:
- RPM(Requests Per Minute):单位时间请求数量。
- TPM(Tokens Per Minute):单位时间输入和输出 Token 数。
- 并发限制:同时执行的请求数量。
只控制 RPM 不够。例如每分钟 60 个请求的配额,如果每个请求包含 100k Token,很快就会耗尽 TPM;只控制并发也不够,因为少量超长请求可能消耗全部 Token 预算。
对每个请求,可以估算其资源需求:
其中第一项代表一个请求, 是预估输入 Token 加最大输出 Token,第三项代表一个并发槽位。只有三个资源都可用时,才应放行请求。
2. 令牌桶模型
令牌桶用来限制平均速率并允许有限突发。设:
- 桶容量为 ;
- 令牌补充速率为 个令牌/秒;
- 当前令牌数为 ;
- 请求消耗 个令牌。
经过 秒后:
当 时,请求被允许,并更新:
否则请求等待,或直接返回限流错误。
请求数限流可以让每个请求消耗一个令牌;Token 限流则让请求消耗预估 Token 数。生产系统往往需要两个桶:
请求桶:每次请求消耗 1
Token 桶:每次请求消耗 estimated_input + max_output
如果估算只使用实际输入 Token,而不为最大输出预留空间,多个请求可能同时通过本地检查,却在输出阶段触发供应商 TPM 限制。
3. 并发信号量解决的是不同问题
信号量限制正在执行的请求数:
import asyncio
concurrency = asyncio.Semaphore(32)
async def call_model(client, **kwargs):
async with concurrency:
return await client.responses.create(**kwargs)
信号量不限制每分钟速率,也不限制 Token 数。一个请求可能占用信号量几秒,却消耗几十万 Token;另一个请求可能只需几十 Token。实际系统通常组合:
请求进入
├── 检查用户/租户配额
├── 估算 Token 并申请 Token 桶
├── 申请并发信号量
├── 调用上游
└── 释放并发槽位,按实际用量修正配额
4. 预估 Token 与实际 Token 的误差
输入 Token 通常可以在发送前通过对应模型的 tokenizer 估算,但不同模型的 tokenizer 不同,聊天消息还包含角色、字段和模板开销。输出 Token 在请求前只能按 max_output_tokens 或历史分布估计。
可采用保守估计:
其中 是根据历史数据设定的平均输出比例。如果系统目标是严格不超过供应商 TPM,应该按最大值预留;如果目标是提高吞吐量,可以使用保守的历史分位数,但要监测估算误差。
请求结束后,使用响应中的实际用量修正成本和限流统计。响应没有 usage 字段时,不能把估算值伪装成精确值,应标记为 estimated。
六、重试、限流和降级的联合故障路径
一个生产请求的典型控制路径如下:
flowchart TD
A[收到业务请求] --> B{预算和权限检查}
B -- 失败 --> Z1[拒绝或返回降级结果]
B -- 通过 --> C[估算输入/输出 Token]
C --> D{本地限流与并发许可}
D -- 无许可 --> Z2[排队或快速失败]
D -- 有许可 --> E[发送上游请求]
E --> F{响应结果}
F -- 成功 --> G[解析、记录 usage、返回]
F -- 429/暂时性错误 --> H{重试预算仍足够?}
H -- 是 --> I[遵守 Retry-After/指数退避]
I --> E
H -- 否 --> Z3[切换模型、供应商或返回错误]
F -- 参数/权限/策略错误 --> Z4[不重试,记录原因]
E --> J{用户取消或截止时间到?}
J -- 是 --> K[关闭流、取消任务、记录取消]
这里的“降级”不是任意换一个模型。模型切换必须考虑:
- 输入上下文长度是否支持;
- 工具调用和结构化输出能力是否兼容;
- 输出语言和质量是否满足业务要求;
- 数据权限是否允许发送到另一供应商;
- 成本和延迟是否符合预算。
例如,支持 JSON Schema 的模型切换到只支持文本 JSON 的模型后,原先的解析保证可能消失。兼容层必须把这种能力差异显式暴露给上层,而不是静默宣称“都支持”。
七、兼容层:统一调用接口,不统一模型语义
1. 协议兼容不等于能力兼容
许多自建或第三方服务提供 OpenAI-compatible API,意味着它们使用相似的路径、请求字段或响应格式。这可以降低迁移成本,但不能推出以下结论:
- 参数含义完全相同;
- Token 计算方式相同;
- 工具调用格式相同;
- 流式事件类型相同;
- 取消行为相同;
- 错误码和限流头相同;
- 结构化输出具有相同约束强度;
- 同名模型具有相同质量和上下文长度。
例如,两个服务都接受 temperature,并不代表它们对该参数的范围、默认值和采样实现相同。某些推理模型可能忽略该参数,某些服务可能直接返回 400。
2. 兼容层应划分为四层
一个可维护的适配器通常包含:
第一层:统一请求模型
应用内部不要把供应商 SDK 类型直接传播到业务代码。定义自己的请求对象:
from dataclasses import dataclass
from typing import Any, AsyncIterator, Protocol
@dataclass
class GenerateRequest:
model: str
messages: list[dict[str, Any]]
max_output_tokens: int | None = None
temperature: float | None = None
tools: list[dict[str, Any]] | None = None
response_format: dict[str, Any] | None = None
这里的字段代表业务需要,不代表所有供应商都能原样支持。
第二层:能力声明
@dataclass(frozen=True)
class ProviderCapabilities:
streaming: bool
cancellation: bool
tools: bool
structured_output: bool
usage_in_stream: bool
max_context_tokens: int | None
调用前根据能力决定策略。例如:
- 没有原生结构化输出时,可以要求文本 JSON,再进行严格校验和有限修复;
- 没有流式能力时,只能等待完整响应;
- 没有取消能力时,至少关闭客户端连接并停止本地转发;
- 没有流式 usage 时,成本统计可能只能在结束后获得。
第三层:供应商映射
class LLMProvider(Protocol):
async def generate(self, request: GenerateRequest):
...
async def stream(self, request: GenerateRequest) -> AsyncIterator[Any]:
...
每个供应商适配器负责把内部请求转换为供应商请求,把供应商响应转换为内部事件。业务代码不应出现:
if provider == "xxx":
...
elif provider == "yyy":
...
这种分支应集中在适配器中。
第四层:统一错误和观测
内部错误至少应包含:
@dataclass
class NormalizedError:
category: str # auth, invalid_request, rate_limit, timeout, upstream, policy
retryable: bool
provider: str
status_code: int | None
request_id: str | None
raw_message: str
上层只依据 category 和 retryable 决策;日志中保留供应商原始错误和请求 ID,便于排查。
3. 兼容层的降级规则必须显式
下面是一个合理的能力决策示例:
def choose_mode(cap: ProviderCapabilities, need_stream: bool) -> str:
if need_stream and cap.streaming:
return "stream"
if need_stream and not cap.streaming:
# 不能伪造实时流,只能使用完整响应后模拟一次发送
return "buffer_then_emit"
return "complete"
buffer_then_emit 只是用户界面层面的兼容,不是真正的流式输出。它会隐藏首 Token 延迟,并且无法让用户提前停止服务端计算。接口名称或响应元数据应让调用方知道这是模拟流。
4. 参数映射不能无条件透传
假设内部请求有:
temperature=0.2
max_output_tokens=1000
适配器不能简单地把所有字段 **kwargs 传给每个服务。应逐字段检查:
payload = {
"model": req.model,
"messages": req.messages,
}
if req.max_output_tokens is not None:
payload["max_output_tokens"] = req.max_output_tokens
if req.temperature is not None and provider_supports_temperature:
payload["temperature"] = req.temperature
对于不支持的字段,有三种合法策略:
- 明确拒绝并返回配置错误;
- 记录警告后忽略;
- 转换为供应商等价字段。
默认静默忽略会造成最难诊断的问题:应用以为参数生效,模型实际没有使用。
八、端到端服务示例:把流式、取消和观测连接起来
下面使用 FastAPI 展示一个简化的 SSE 转发服务。它的重点是生命周期和取消传播,不是生产级认证或完整供应商适配。
import os
from contextlib import asynccontextmanager
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
from openai import AsyncOpenAI
client: AsyncOpenAI | None = None
@asynccontextmanager
async def lifespan(app: FastAPI):
global client
client = AsyncOpenAI(
api_key=os.environ["OPENAI_API_KEY"],
timeout=120.0,
max_retries=0,
)
yield
await client.close()
app = FastAPI(lifespan=lifespan)
@app.get("/chat/stream")
async def chat_stream(request: Request, q: str):
assert client is not None
async def event_generator():
stream = await client.responses.create(
model=os.environ["OPENAI_MODEL"],
input=q,
stream=True,
)
try:
async for event in stream:
# 浏览器断开后,is_disconnected 通常会变为 True
if await request.is_disconnected():
break
if event.type == "response.output_text.delta":
text = event.delta.replace("\n", "\\n")
yield f"data: {text}\n\n"
elif event.type == "response.completed":
yield "event: done\ndata: {}\n\n"
finally:
# 具体 SDK 是否支持显式关闭,应查对应版本文档。
# 即使没有公开 close 方法,也必须让生成器退出并释放请求上下文。
pass
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
},
)
请求:
curl -N "http://localhost:8000/chat/stream?q=解释HTTP流式响应"
-N 会让 curl 尽量不缓冲输出。预期结果类似:
data: HTTP 流式响应...
data: ...
event: done
data: {}
示例中的 X-Accel-Buffering: no 只对某些 Nginx 配置有帮助,不能替代网关层配置。还需要检查:
- 反向代理是否启用了响应缓冲;
- 是否压缩了 SSE 响应;
- 是否设置了足够的空闲超时;
- 是否正确转发客户端断开事件;
- 应用服务器是否在生成器退出时释放上游连接。
生产实现还应为每个请求生成 trace_id,并记录:
request_id
provider_request_id
model
tenant_id
input_tokens
output_tokens
time_to_first_token
total_latency
retry_count
cancel_reason
finish_reason
estimated_cost
其中 TTFT(Time To First Token) 是从请求发出到收到第一个有效文本增量的时间。TTFT 高可能来自排队、连接建立、模型 Prefill 或供应商调度;不能仅凭总耗时判断。
九、成本和 Token 统计必须与生命周期绑定
LLM 成本通常可表示为:
其中:
- 是输入 Token;
- 是输出 Token;
- 、 是供应商定价;
- 是外部工具或检索成本;
- 表示缓存命中、缓存写入或缓存读取相关费用,具体是否存在取决于供应商。
流式与非流式的计费通常由供应商服务端的实际使用量决定,而不是由客户端是否打印了文本决定。用户在第一个 Token 后取消,服务端可能已经生成了一部分输出,因此不能把“用户没有看到完整答案”解释成“没有输出成本”。
如果响应没有最终 usage,需要分别记录:
usage.input_tokens = actual | estimated | unavailable
usage.output_tokens = actual | estimated | unavailable
不要把 max_output_tokens 当成实际输出 Token。前者是上限,后者是实际生成量。
缓存也会影响调用逻辑。若缓存命中,系统可能不会产生一次新的模型请求;若缓存键只包含用户问题而不包含系统指令、模型版本、工具定义和权限范围,可能返回不适用甚至越权的结果。一个安全缓存键通常至少考虑:
hash(
normalized_system_prompt,
user_input,
model,
model_parameters,
tool_schema,
retrieval_context_version,
authorization_scope
)
十、与 Prompt 工程和结构化输出的边界
客户端层不能修复错误的 Prompt,也不能把一个不支持严格结构化输出的模型变成严格校验器。
一次结构化生成至少包含三个不同约束:
- 指令约束:Prompt 要求输出某种格式。
- 协议约束:API 是否支持 JSON、Schema 或工具调用。
- 验证约束:客户端是否用 Pydantic、JSON Schema 等再次验证。
例如,下面的响应即使是合法 JSON,也可能不满足业务要求:
{"name": "Alice", "age": "unknown"}
如果 age 必须是整数,客户端应在模型响应后执行严格校验:
from pydantic import BaseModel, ValidationError
class Person(BaseModel):
name: str
age: int
try:
person = Person.model_validate_json(response.output_text)
except ValidationError:
# 记录原始输出,决定是否进行一次受控修复或返回错误
raise
“解析失败就无限重试”是错误做法。模型输出校验失败时,最多进行少量、明确目的的修复请求,并且要重新计算 Token、延迟和预算。若连续失败,应返回可诊断错误,而不是继续消耗资源。
十一、常见错误和诊断方式
1. 把 429 一律快速重试
表现:
请求量上升 → 429 增多 → 客户端立即重试 → 上游压力更高 → 429 更多
诊断方法:
- 记录每次尝试的 HTTP 状态码;
- 记录是否存在
Retry-After; - 统计原始请求数与实际上游尝试数;
- 查看 Token 限额是否先于请求数耗尽。
修复方式是令牌桶、指数退避、重试上限和排队策略的组合,而不是增加重试次数。
2. 只设置总超时,不设置流空闲超时
表现是部分请求长期占用连接,应用并发槽位逐渐耗尽。诊断时查看:
active_streams
last_event_age
request_total_duration
client_disconnect_count
如果 last_event_age 持续超过阈值,应主动取消;但阈值应结合模型、上下文长度和供应商调度特征设置。
3. 认为关闭浏览器就会停止模型生成
表现是用户已经离开页面,但供应商 usage 仍增加。原因是浏览器断开、应用任务取消、上游模型停止之间存在多个独立边界。应在应用层确认:
- 是否检测到客户端断开;
- 是否取消了本地任务;
- 是否关闭了上游响应;
- 供应商是否提供服务端取消接口;
- 取消后的实际 Token 和成本如何统计。
4. 把流事件当作完整文本
表现包括乱码、重复文本、JSON 解析失败和工具参数截断。正确做法是按事件类型处理,并在完整输出或明确完成事件后再执行最终解析。
5. 直接依赖供应商 SDK 类型
表现是更换模型供应商时,业务层到处出现不同的消息格式、异常类型和流事件判断。适配器应在边界处完成转换,业务层使用自己的请求、事件和错误模型。
6. 让 SDK 和应用层同时重试
表现是配置了 SDK 的默认重试,又在外层包装一个重试循环,最终一次业务请求产生过多上游尝试。应明确重试所有权:
SDK 负责重试,应用只处理最终错误
或
SDK 不重试,应用统一负责分类、退避、预算和观测
后一种更容易与租户配额、成本预算和多供应商降级结合,但需要自己正确处理 SDK 异常。
十二、生产设计的最小检查表
在上线一个 LLM API 客户端前,至少应验证以下行为:
客户端和连接
- 客户端是否复用,而不是每次请求创建?
- 应用退出时是否关闭连接池?
- 同步和异步调用是否没有相互阻塞?
- 连接、读取、总耗时是否有明确预算?
流式输出
- 是否按事件类型解析,而不是按 Token 假设?
- 是否处理完成、失败、工具调用和拒答事件?
- 转发层是否关闭缓冲?
- 下游变慢时是否存在背压或有界队列?
取消
- 用户断开是否传播到本地异步任务?
- 取消后是否释放信号量和响应体?
- 是否区分用户取消、超时和上游失败?
- 是否了解供应商对服务端取消和计费的保证边界?
重试
- 是否只重试明确的暂时性错误?
- 是否使用指数退避和抖动?
- 是否遵守
Retry-After? - 工具调用和外部副作用是否有幂等键?
- 是否限制总尝试次数和重试并发?
限流
- 是否分别控制 RPM、TPM 和并发?
- Token 估算是否考虑输入、最大输出和消息模板?
- 429 后是否降低压力,而不是立即重试?
- 是否有租户级、模型级和供应商级预算?
兼容层
- 是否有内部统一请求和响应模型?
- 是否显式声明供应商能力?
- 不支持的参数是否拒绝、转换或明确忽略?
- 是否保留原始供应商请求 ID 和错误?
- 是否验证结构化输出,而不是只检查 HTTP 200?
结语
LLM API 的核心难点不在于发送一个 JSON,而在于管理一次远程生成任务的完整生命周期。客户端负责连接和资源;流式输出负责增量事件与背压;取消负责传播“不再需要结果”的信号;重试负责在不放大故障的前提下恢复暂时性错误;限流负责同时约束请求、Token 和并发资源;兼容层负责隔离协议差异,并明确暴露能力差异。
真正稳定的实现不会把这些问题分别处理成几个独立工具,而会把它们连接在同一条控制链路中:
权限与预算
→ Token 估算
→ 本地限流
→ 上游调用
→ 流式/完整响应解析
→ 取消或超时传播
→ 有边界的重试
→ usage、Trace、成本和错误记录
当模型、数据、权限、评测和成本都被视为同一个生产系统的组成部分时,LLM API 客户端才不只是一个 SDK 包装器,而是生成式 AI 系统中的资源调度、故障隔离和语义兼容边界。
系列导航与关联阅读
- 系列入口:AI 工程完整学习路线:从机器学习与 Transformer 到 RAG、Agent 和生产治理
- 上一篇:Prompt 工程:指令层级、上下文、Few-shot、结构化输出和测试
- 下一篇:LLM 结构化输出与工具调用:Schema、循环、幂等、授权和确认
- 延伸:AI 可观测性与成本治理:Trace、Token、TTFT、预算、缓存和降级
官方资料
本文依据研究论文、标准组织与主流框架官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论