繁中 ▾

轉接 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 調到上限 32,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())

若多個進程或機器共用同一把 API 金鑰,記憶體節拍器就不夠用了,因為速率限制是按金鑰合併統計的。此時需共用令牌桶,常見做法是用 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 金鑰