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 请求通常经历以下阶段:

  1. 构造输入消息、工具定义、结构化输出约束和模型参数。
  2. 从客户端连接池获取或建立 TCP/TLS 连接。
  3. 发送 HTTP 请求。
  4. 等待服务端接受请求并开始推理。
  5. 接收完整响应,或持续接收流式事件。
  6. 解析响应、记录 Token 和耗时。
  7. 释放响应体、连接和本地任务资源。

因此,客户端对象通常应在进程或应用生命周期内复用,而不是每个请求都重新创建。复用客户端可以复用连接池、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 超时。模型可能在生成较长答案时,连续事件之间暂时存在间隔;如果空闲超时过短,会把正常生成误判为故障。反过来,如果没有空闲超时,半断开的连接可能长期占用并发槽位。

工程上可以使用两个时间约束:

TtotalTbudgetT_{\text{total}} \leq T_{\text{budget}}

ΔteventTidle\Delta t_{\text{event}} \leq T_{\text{idle}}

其中 TtotalT_{\text{total}} 是请求总预算,Δtevent\Delta t_{\text{event}} 是相邻流事件之间的时间间隔,TidleT_{\text{idle}} 是允许的最大空闲时间。两者分别防止“请求整体过慢”和“连接失活但不退出”。


二、流式输出:传输层、事件层和文本层不是一回事

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”,原因包括:

  1. 服务端可能按内部批次发送多个 Token。
  2. 一个 Token 可能对应一个 UTF-8 字节序列的一部分或多个字符。
  3. 事件可能承载文本增量、工具调用参数、状态变化、错误和完成信息。
  4. 代理、网关或 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_KEYOPENAI_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. 取消有多个层级

“取消请求”至少可能指:

  1. 业务层取消:用户点击停止、页面关闭或上游任务失效。
  2. 应用任务取消:Python 的 asyncio.Task 被取消。
  3. HTTP 连接取消:客户端关闭响应流或断开连接。
  4. 供应商任务取消:通过供应商提供的取消接口通知服务端停止生成。
  5. 模型计算取消:推理服务器真正中止 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. 取消的幂等性

取消操作通常应设计为幂等:

cancel(cancel(x))=cancel(x)cancel(cancel(x)) = cancel(x)

用户可能重复点击停止,网络重试也可能重复发送取消请求。服务端收到第一次取消后,第二次应返回“已完成”“已取消”或“资源不存在”等可接受状态,而不应破坏其他请求。

应用层还应记录取消原因:

  • user_cancelled
  • client_disconnected
  • deadline_exceeded
  • upstream_failed
  • budget_exceeded

这些原因对成本分析和用户体验诊断不同。把所有中止都记成“模型错误”会掩盖真实问题。


四、重试:先判断故障是否可重试,再决定等待多久

1. 错误分类比重试代码更重要

常见错误可按可重试性分组:

错误类型 典型原因 通常是否重试
DNS、连接失败 网络或服务不可达 可以有限重试
连接/读取超时 网络抖动或服务过慢 视请求是否已到达服务端
HTTP 408 请求超时 通常可重试
HTTP 429 限流或配额不足 按服务端提示延迟
HTTP 500、502、503、504 服务端暂时故障 可以有限重试
HTTP 400 参数、格式或上下文错误 不应原样重试
HTTP 401、403 认证或权限错误 不应重试
内容安全拒绝 输入或输出违反策略 通常不应原样重试
业务校验失败 模型输出不符合格式 应修复提示或降级,不是 HTTP 重试

HTTP 状态码只是第一层信号。响应体中的错误类型、请求 ID 和供应商特定错误码往往更准确。

2. 流式请求的重试更危险

非流式请求在客户端没有收到响应时,无法确定服务端是否已经执行了请求。例如:

客户端发送成功
服务端完成模型生成
服务端准备返回
连接在返回前断开
客户端看到超时

此时客户端重试,可能导致同一个请求执行两次。对于纯文本生成,重复通常只是浪费成本;对于模型调用工具、发送邮件、创建订单等场景,重复可能造成真实副作用。

所以需要区分:

  • 生成操作:通常无外部副作用,但可能重复计费。
  • 工具执行:可能产生外部副作用,必须使用幂等键或业务去重。
  • 流式输出:已经向用户展示了部分结果后再重试,不能简单拼接两次答案。

如果流式连接在输出中途断开,应用有三种策略:

  1. 标记失败,要求用户重新生成;
  2. 使用相同请求标识重新请求,并覆盖旧输出;
  3. 继续请求并尝试恢复,但需要服务端支持可恢复游标或响应 ID。

第三种不能通过“把已输出文本放入新 Prompt”简单实现,因为模型可能重新生成不同内容,也会增加输入 Token。

3. 指数退避和抖动

常见退避时间为:

dk=min(dmax,d02k)+U(0,j)d_k = \min(d_{\max}, d_0 \cdot 2^k) + U(0, j)

其中:

  • kk 是第 kk 次重试;
  • d0d_0 是初始等待时间;
  • dmaxd_{\max} 是最大等待时间;
  • U(0,j)U(0,j) 是范围为 00jj 的随机抖动。

随机抖动可以避免大量客户端在同一时刻同时重试,形成“惊群”。

若响应包含 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 倍。重试应受以下约束:

  • 单请求最大尝试次数;
  • 单请求总时间预算;
  • 服务实例级重试并发上限;
  • 全局错误率熔断;
  • 供应商配额和成本预算。

设原始请求速率为 λ\lambda,平均尝试次数为 aa,则上游实际请求速率近似为:

λupstream=λa\lambda_{\text{upstream}} = \lambda \cdot a

当 429 或 503 增多时,aa 可能上升,进一步推高 λupstream\lambda_{\text{upstream}},形成正反馈。限流和重试必须一起设计,否则“重试恢复故障”会变成“重试制造故障”。


五、限流:请求数、并发数和 Token 数是三个不同资源

1. RPM、TPM 和并发限制

LLM 服务常见三类限制:

  • RPM(Requests Per Minute):单位时间请求数量。
  • TPM(Tokens Per Minute):单位时间输入和输出 Token 数。
  • 并发限制:同时执行的请求数量。

只控制 RPM 不够。例如每分钟 60 个请求的配额,如果每个请求包含 100k Token,很快就会耗尽 TPM;只控制并发也不够,因为少量超长请求可能消耗全部 Token 预算。

对每个请求,可以估算其资源需求:

qi=(1, t^i, 1)q_i = (1,\ \hat{t}_i,\ 1)

其中第一项代表一个请求,t^i\hat{t}_i 是预估输入 Token 加最大输出 Token,第三项代表一个并发槽位。只有三个资源都可用时,才应放行请求。

2. 令牌桶模型

令牌桶用来限制平均速率并允许有限突发。设:

  • 桶容量为 BB
  • 令牌补充速率为 rr 个令牌/秒;
  • 当前令牌数为 bb
  • 请求消耗 cc 个令牌。

经过 Δt\Delta t 秒后:

b=min(B, b+rΔt)b' = \min(B,\ b + r\Delta t)

bcb' \geq c 时,请求被允许,并更新:

bafter=bcb_{\text{after}} = b' - c

否则请求等待,或直接返回限流错误。

请求数限流可以让每个请求消耗一个令牌;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 或历史分布估计。

可采用保守估计:

t^=tinput+αtmax_output\hat{t} = t_{\text{input}} + \alpha \cdot t_{\text{max\_output}}

其中 0<α10 < \alpha \leq 1 是根据历史数据设定的平均输出比例。如果系统目标是严格不超过供应商 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

上层只依据 categoryretryable 决策;日志中保留供应商原始错误和请求 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

对于不支持的字段,有三种合法策略:

  1. 明确拒绝并返回配置错误;
  2. 记录警告后忽略;
  3. 转换为供应商等价字段。

默认静默忽略会造成最难诊断的问题:应用以为参数生效,模型实际没有使用。


八、端到端服务示例:把流式、取消和观测连接起来

下面使用 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 成本通常可表示为:

C=pinTin106+poutTout106+Ctool+CcacheC = p_{\text{in}} \cdot \frac{T_{\text{in}}}{10^6} + p_{\text{out}} \cdot \frac{T_{\text{out}}}{10^6} + C_{\text{tool}} + C_{\text{cache}}

其中:

  • TinT_{\text{in}} 是输入 Token;
  • ToutT_{\text{out}} 是输出 Token;
  • pinp_{\text{in}}poutp_{\text{out}} 是供应商定价;
  • CtoolC_{\text{tool}} 是外部工具或检索成本;
  • CcacheC_{\text{cache}} 表示缓存命中、缓存写入或缓存读取相关费用,具体是否存在取决于供应商。

流式与非流式的计费通常由供应商服务端的实际使用量决定,而不是由客户端是否打印了文本决定。用户在第一个 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,也不能把一个不支持严格结构化输出的模型变成严格校验器。

一次结构化生成至少包含三个不同约束:

  1. 指令约束:Prompt 要求输出某种格式。
  2. 协议约束:API 是否支持 JSON、Schema 或工具调用。
  3. 验证约束:客户端是否用 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 仍增加。原因是浏览器断开、应用任务取消、上游模型停止之间存在多个独立边界。应在应用层确认:

  1. 是否检测到客户端断开;
  2. 是否取消了本地任务;
  3. 是否关闭了上游响应;
  4. 供应商是否提供服务端取消接口;
  5. 取消后的实际 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 系统中的资源调度、故障隔离和语义兼容边界。


系列导航与关联阅读

官方资料

本文依据研究论文、标准组织与主流框架官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。