中转 API 稳定性实践:超时、重试、限流队列与流式断连

把中转 API 接进生产之后,真正消耗你精力的往往不是接入,而是那些偶发的异常:几十次请求里冒出一两个超时,批处理跑到一半被限流,流式输出在中途断掉。这篇文章用运维的方式一项项处理:先把错误分类,再给出可运行的退避重试、并发队列、断流续写和用量监控代码。

更新于

要点

  1. 先分类再重试:429 和 503 值得退避重试,400、401、402、403 重试只会浪费时间。
  2. 每分钟 300 次的上限,要用“信号量控并发 + 间隔控速率”两层来守住,只用其中一层不够。
  3. 流式请求要把读超时当成“数据块间隔”来设置,中断后保留已收到的内容再续写。
  4. 把每次请求的 usage 落盘,就能提前发现成本异常和提示词膨胀。

先把错误分成三类

稳定性问题的第一步不是加重试,而是知道什么该重试。按处理方式可以把常见失败分成三类:

类别典型表现处理方式
瞬时可恢复429 限流、503 upstream_busy、连接被重置、超时指数退避后重试,限制总次数
请求本身有问题400(输入加 max_tokens 超过 100,000、请求体过大)、403 content_blocked不重试,修正请求或直接反馈给用户
账户状态问题401 密钥无效、402 no_credit不重试,告警给负责人,充值或更换密钥

这张表最好写进代码注释里。一个常见的事故是:余额用完后接口开始稳定返回 402,而你的重试逻辑没有区分状态码,于是每个请求都重试五次,白白放大了五倍的流量,同时把日志刷成一片红。

超时要分层设置

只设一个总超时是不够的。建议把超时拆成三层来考虑:

  • 连接超时:建立 TCP 和 TLS 连接的时间,通常 5 到 10 秒足够,连不上就应该尽快失败。
  • 读超时:等待响应数据的时间。非流式请求要等整段文字生成完才返回,所以它必须覆盖最长输出;流式请求则变成“两个数据块之间的最长静默”,10 到 30 秒更合理。
  • 业务总超时:你自己的业务能容忍的最长等待,用 asyncio.wait_for 或网关层来保证,包含所有重试在内。

输出越长,非流式请求的耗时越久。如果你把 max_tokens 调到上限 16,000,又只设了 30 秒总超时,必然会在长文场景下频繁超时。更稳妥的做法是长输出一律使用流式,既能给用户即时反馈,也让超时的含义清晰。SDK 里的 timeout 参数可以传一个数字,也可以传分项的配置,按你使用的版本查阅文档。

429 与 503 的指数退避重试

退避的要点有三个:延迟随失败次数翻倍增长、设置上限、加入随机抖动。没有抖动时,一批同时失败的任务会在同一时刻一起重试,把刚恢复的服务再次压垮。下面的函数关闭了 SDK 自带的重试,由我们自己统一控制,并且只对 429 和 503 以及连接类错误重试:

import asyncio
import os
import random

from openai import AsyncOpenAI, APIConnectionError, APIStatusError, APITimeoutError

client = AsyncOpenAI(
    base_url="https://api.llmzhongzhuan.com/v1",
    api_key=os.environ["API_KEY"],
    timeout=60.0,
    max_retries=0,  # 重试由下面的函数统一控制
)

RETRY_STATUS = {429, 503}


async def chat_with_retry(messages, max_attempts=5, base_delay=1.0, cap=20.0):
    for attempt in range(1, max_attempts + 1):
        try:
            resp = await client.chat.completions.create(
                model="uncensored", messages=messages, max_tokens=800
            )
            return resp
        except APIStatusError as e:
            if e.status_code not in RETRY_STATUS or attempt == max_attempts:
                raise  # 400/401/402/403 重试没有意义
            reason = f"HTTP {e.status_code}"
        except (APIConnectionError, APITimeoutError) as e:
            if attempt == max_attempts:
                raise
            reason = type(e).__name__
        delay = min(cap, base_delay * 2 ** (attempt - 1))
        delay = random.uniform(0, delay)  # 全抖动,避免所有任务同时醒来
        print(f"第 {attempt} 次失败({reason}),{delay:.1f}s 后重试")
        await asyncio.sleep(delay)

503 的含义是服务端暂时繁忙,建议几秒后重试,上面的基准延迟 1 秒、上限 20 秒、最多 5 次,足以覆盖这种场景。对于 429,如果你发现它频繁出现,说明问题不在重试,而在发送端没有做速率控制,请看下一节。

每分钟 300 次限流下的并发队列

每个密钥每分钟 300 次请求,折算下来平均每秒 5 次。批处理任务最容易在这里踩坑:用 asyncio.gather 一次性发出上千个请求,前几十个瞬间用光额度,其余全部收到 429。

可靠的做法是两层控制。信号量限制“同时在途”的数量,防止本机连接和内存被撑爆;节拍器限制“发出速率”,把请求均匀铺开。两者缺一不可:只有信号量时,如果单个请求很快,每分钟仍可能超过 300 次;只有节拍器时,如果请求很慢,在途请求会越堆越多。示例里把速率定在每分钟 270 次,给多实例、重试和时钟误差留出 10% 的余量:

import asyncio
import os
import time

from openai import AsyncOpenAI

client = AsyncOpenAI(
    base_url="https://api.llmzhongzhuan.com/v1",
    api_key=os.environ["API_KEY"],
    timeout=60.0,
    max_retries=0,
)

MAX_IN_FLIGHT = 8          # 同时在途的请求数
PER_MINUTE = 270           # 留 10% 余量,低于每分钟 300 次的上限


class Pacer:
    """把请求的发出时刻均匀铺开:每 60/PER_MINUTE 秒放行一个。"""

    def __init__(self, per_minute):
        self.interval = 60.0 / per_minute
        self.next_at = 0.0
        self.lock = asyncio.Lock()

    async def wait(self):
        async with self.lock:
            now = time.monotonic()
            delay = self.next_at - now
            if delay > 0:
                await asyncio.sleep(delay)
            self.next_at = max(now, self.next_at) + self.interval


sem = asyncio.Semaphore(MAX_IN_FLIGHT)
pacer = Pacer(PER_MINUTE)


async def summarize(idx, text):
    async with sem:            # 控制并发
        await pacer.wait()     # 控制速率
        resp = await client.chat.completions.create(
            model="uncensored",
            messages=[{"role": "user", "content": f"用两句话概括:{text}"}],
            max_tokens=120,
        )
        return idx, resp.choices[0].message.content


async def main():
    texts = [f"第 {i} 条工单的正文……" for i in range(1000)]
    tasks = [summarize(i, t) for i, t in enumerate(texts)]
    done = 0
    for coro in asyncio.as_completed(tasks):
        idx, out = await coro
        done += 1
        if done % 100 == 0:
            print(f"已完成 {done} 条")


asyncio.run(main())

如果有多个进程或多台机器共用同一把密钥,上面的内存级节拍器就不够了,限流是按密钥合并统计的。此时需要一个共享的令牌桶,常见做法是用 Redis 做计数,或者干脆让所有批处理请求都经过一个单独的队列消费服务,由它统一出口。

流式输出中断后怎么办

流式连接比普通请求更脆弱:它持续的时间长,经过的网络设备多,任何一个环节的空闲超时都可能把它切断。处理原则是“已经收到的内容就是资产”,不要因为中断就全部丢弃。

下面的实现把已收到的片段累积起来,捕获连接和超时异常后,把已有文本当作 assistant 消息再请求一次,让模型接着往下写。这是一种尽力而为的补救,续写的衔接不一定完美,适合长文和对话,不适合要求严格结构的输出,例如 JSON。结构化输出一旦中断,更稳妥的做法是整体重来。注意流式结束时会自动附带一个 usage 分片,代码里同时把它取了出来:

import asyncio
import os

from openai import AsyncOpenAI, APIConnectionError, APITimeoutError

client = AsyncOpenAI(
    base_url="https://api.llmzhongzhuan.com/v1",
    api_key=os.environ["API_KEY"],
    timeout=30.0,   # 流式场景下,这是相邻数据块之间允许的最长静默
    max_retries=0,
)


async def stream_once(messages):
    parts, usage = [], None
    stream = await client.chat.completions.create(
        model="uncensored", messages=messages, max_tokens=600, stream=True
    )
    try:
        async for chunk in stream:
            if chunk.usage:                      # 最后一个分片携带 usage
                usage = chunk.usage
            if chunk.choices and chunk.choices[0].delta.content:
                parts.append(chunk.choices[0].delta.content)
    except (APIConnectionError, APITimeoutError):
        return "".join(parts), None, False       # 中断:返回已收到的部分
    return "".join(parts), usage, True


async def stream_with_resume(question, max_resume=2):
    messages = [{"role": "user", "content": question}]
    text = ""
    for _ in range(max_resume + 1):
        part, usage, finished = await stream_once(messages)
        text += part
        if finished:
            return text, usage
        # 把已生成部分作为 assistant 消息,请模型接着写
        messages = [
            {"role": "user", "content": question},
            {"role": "assistant", "content": text},
            {"role": "user", "content": "请从上一条回复的结尾处继续,不要重复已写的内容。"},
        ]
    return text, None


if __name__ == "__main__":
    out, usage = asyncio.run(stream_with_resume("用三段话说明为什么日志要带 request id。"))
    print(out)
    print("usage:", usage)

续写时要记得,已发送的文本也会算作输入 token,所以每次续写都会有一点额外成本,限制 max_resume 的次数是必要的。

降级、熔断与上线前自查

重试解决的是瞬时故障,如果故障持续几分钟,继续重试只会堆积请求。建议在重试之上再加一层熔断:连续失败达到阈值,比如一分钟内失败十次,就暂停向外发请求一小段时间,期间直接走降级逻辑,例如返回缓存结果、提示用户稍后再试,或者把任务放回队列延后处理。冷却时间结束后只放行少量探测请求,成功了再恢复全量。

降级也要提前想清楚。面向用户的聊天功能可以返回一句友好的忙碌提示;离线的批处理任务则应该把失败条目写进失败表,等整体跑完后再集中补跑,而不是在主流程里无限等待。无论哪种,都要保证失败不会静默吞掉,至少留下一行带 request id 的日志。

最后给一份上线前的自查清单:是否关闭了 SDK 的隐式重试并只保留自己的一处;是否对 400、401、402、403 直接放弃;总超时是否把所有重试都算了进去;批处理是否有速率控制而不是一次性并发;流式是否处理了中断;usage 是否落盘;余额是否有告警。逐项通过后,再去压测,而不是反过来。

用 usage 做监控和对账

每次响应都包含 usage,流式则出现在最后一个分片。把它连同功能名一起落盘,是成本最便宜的监控手段。本站输入价格为每百万 token 0.25 美元、输出为 1.00 美元,可以在本地直接估算每次请求的费用:

import csv
import time

LOG = "usage_log.csv"


def record_usage(feature, resp):
    u = resp.usage
    cost = u.prompt_tokens * 0.25 / 1_000_000 + u.completion_tokens * 1.00 / 1_000_000
    with open(LOG, "a", newline="", encoding="utf-8") as f:
        csv.writer(f).writerow(
            [int(time.time()), feature, u.prompt_tokens, u.completion_tokens, f"{cost:.6f}"]
        )

落盘之后,建议每天看三个指标:各功能的输入 token 中位数,用来发现提示词是否悄悄膨胀;输出 token 占比,输出单价是输入的四倍,所以它通常是账单的大头;以及 429 与 503 的比例,如果比例上升,说明发送速率或重试策略需要调整。再加一条余额告警,在 402 出现之前就通知负责人充值。更多接入层面的细节可以参考 框架配置,常见问题则汇总在 常见问题 里。

常见问题

收到 429 应该立刻重试吗?

不应该。先退避等待并加随机抖动,同时检查发送端是否超过每分钟 300 次。长期出现 429 说明需要加节拍器,而不是加重试次数。

503 upstream_busy 要等多久再试?

几秒即可。建议从 1 秒起指数退避,上限 20 秒左右,并限制总次数,超过后向上层返回明确的失败。

流式请求断了,已经收到的内容要丢掉吗?

不要丢。可以保留已收到的文本,把它作为 assistant 消息再请求续写;若是 JSON 等严格结构,则整体重新请求更稳妥。

怎样确认每次请求花了多少钱?

读取响应里的 usage,用输入 token 乘以每百万 0.25 美元、输出 token 乘以每百万 1.00 美元即可估算,流式时 usage 在最后一个分片。

多台机器共用一把密钥时限流怎么算?

按密钥合并统计,所有机器的请求加起来每分钟不能超过 300 次。需要共享计数或统一出口队列来协调。

只需填写表单即可获取密钥

创建账户,复制密钥,修改 Base URL。配置就是这么简单。

获取 API 密钥