プロキシ API の安定性実践:タイムアウト、リトライ、レート制限キューとストリーミング切断
プロキシ API を本番環境に接続した後、実際にリソースを消費するのは接続そのものではなく、偶発的なエラーです。数十回のリクエスト中にタイムアウトが発生したり、バッチ処理中にレート制限に達したり、ストリーミングが途中で切断されたりします。この記事では、エラーを分類し、実行可能なバックオフ・リトライ、同時リクエストキュー、ストリーミングの再接続、使用量監視のコードを順に解説します。
ポイント
- まず分類してからリトライ:429 と 503 はバックオフしてリトライする価値がありますが、400、401、402、403 はリトライしても時間を無駄にするだけです。
- 毎分 300 回の上限は、「セマフォによる同時実行数の制御」と「インターバルによるレート制御」の 2 層で守る必要があります。片方だけでは不十分です。
- ストリーミングリクエストでは、読み込みタイムアウトを「データブロック間の間隔」として設定し、切断後は受信済みの内容を保持して続きを生成します。
- 各リクエストの usage を保存しておけば、コストの異常やプロンプトの膨張を事前に検知できます。
まずエラーを 3 つのクラスに分類する
安定性問題の第一歩はリトライを増やすことではなく、何をリトライすべきかを知ることです。処理方法に応じて、一般的な失敗を 3 つのクラスに分類できます。
| カテゴリ | 典型的な症状 | 処理方法 |
|---|---|---|
| 一時的な回復可能エラー | 429 レート制限、503 upstream_busy、接続のリセット、タイムアウト | 指数関数的バックオフ後にリトライし、試行回数を制限する |
| リクエスト自体に問題がある | 400(入力と max_tokens の合計が 100,000 を超える、リクエストボディが大きい)、403 content_blocked | リトライせず、リクエストを修正するかユーザーに直接フィードバックする |
| アカウント状態の問題 | 401 無効なキー、402 no_credit | リトライせず、担当者へアラートを送り、チャージまたはキーの交換を行う |
この表はコードコメントに記述するのが最適です。よくある事故は、残高が枯渇した後に API が安定して 402 を返し、リトライロジックがステータスコードを区別しないため、各リクエストが 5 回リトライしてトラフィックが 5 倍に増大し、ログがエラーで埋め尽くされることです。
タイムアウトは階層化して設定する
単一の総タイムアウトのみを設定するのは不十分です。タイムアウトを 3 層に分けて検討することをお勧めします。
- 接続タイムアウト:TCP および TLS 接続を確立する時間。通常 5〜10 秒で十分であり、接続できない場合は速やかに失敗させるべきです。
- 読み込みタイムアウト:レスポンスデータを待つ時間。非ストリーミングリクエストは全文生成完了まで待つ必要があるため、最長の出力をカバーする必要があります。ストリーミングの場合、「2 つのデータチャンク間の最長の無音」になり、10〜30 秒がより適切です。
- ビジネス全体のタイムアウト:ビジネスが許容できる最長の待機時間。
asyncio.wait_forやゲートウェイ層で保証し、すべてのリトライを含めます。
出力が長ければ長いほど、非ストリーミングリクエストの所要時間は長くなります。max_tokens を上限の 32,000 に設定し、総タイムアウトを 30 秒に設定すると、長文生成で頻繁にタイムアウトが発生します。より確実な方法は、長文出力には常にストリーミングを使用し、即時フィードバックを提供するとともに、タイムアウトの定義を明確にすることです。SDK の timeout パラメータには数値または詳細な設定を指定でき、使用しているバージョンのドキュメントを参照してください。
429 と 503 の指数関数的バックオフによるリトライ
バックオフの要点は 3 つです:遅延は失敗回数に応じて倍増、上限を設定、ランダムジッターを追加します。ジッターがないと、一斉に失敗したタスクは同時にリトライし、回復したばかりのサービスを再び圧迫します。以下の関数では 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 を受け取ります。
信頼性の高い方法は 2 層制御です。セマフォは「同時進行中」の数を制限し、ローカル接続とメモリが溢れるのを防ぎます。レートリミッターは「送信レート」を制限し、リクエストを均等に分散します。両方が不可欠です。セマフォのみでは、単一リクエストが非常に速い場合、毎分 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)
続きを書く際は、送信済みのテキストも入力トークンとしてカウントされるため、各続き生成にはわずかな追加コストがかかることを覚えておいてください。max_resume の回数を制限することは必須です。
デグレード、サーキットブレーカーとリリース前の自己チェック
リトライは一時的な障害を解決しますが、障害が数分続いた場合、リトライを続けるとリクエストが蓄積するだけです。リトライの上にサーキットブレーカーを追加することをお勧めします。連続失敗が閾値(例:1 分間で 10 回失敗)に達したら、一時的に外部へのリクエスト送信を停止し、デグレードロジック(キャッシュ結果の返送、ユーザーへの再試行案内、キューへのタスク返却など)に切り替えます。クールダウン期間終了後、少量のプローブリクエストのみを通し、成功すればフルリクエストを再開します。
フェイルバックも事前に明確にしておく必要があります。ユーザー向けチャット機能には親切なビジーメッセージを返し、オフラインのバッチ処理は失敗エントリを失敗テーブルに書き込み、全体の実行後に一括で再実行するようにします。メインフローで無限に待機するのではなく。いずれの場合も、失敗が静かに消えないようにし、少なくとも request id を含むログを残す必要があります。
最後に、運用開始前のチェックリストを提示します:SDK の暗黙的なリトライをオフにし、自前の 1 箇所のみを残したか。400、401、402、403 に対しては即座に諦めたか。総タイムアウトに全てのリトライが含まれているか。バッチ処理にレート制御があり、一度に並列実行していないか。ストリーミングの切断を処理したか。usage を保存したか。残高にアラートが設定されているか。項目を一つずつ通過してから負荷テストを行い、その逆ではないようにします。
usage による監視と照合
各レスポンスには usage が含まれ、ストリーミングの場合は最後のチャンクに現れます。これを機能名とともに保存することは、最もコストの低い監視手法です。当サイトの入力価格は百万トークンあたり 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}"]
)
落盤後、以下の 3 つの指標を毎日確認することをお勧めします:機能別の入力トークンの中央値(プロンプトが静かに肥大化していないか確認するため)、出力トークンの割合(出力単価が入力の 4 倍のため、通常は請求書の大部分を占める)、および 429 と 503 の比率(比率が上昇した場合、送信レートまたはリトライ戦略の調整が必要であることを示します)。さらに残高アラートを追加し、402 が発生する前に担当者へチャージを通知します。接続の詳細については フレームワーク設定 を、よくある質問は よくある質問 にまとめられています。
よくある質問
429を受信したらすぐにリトライすべきか?
すべきではありません。まずバックオフ待機し、ランダムジッターを追加してください。同時に、送信側が毎分300回を超えていないか確認してください。429が長期間発生している場合は、リトライ回数を増やすのではなく、レートリミッターを追加する必要があります。
503 upstream_busyの再試行にはどのくらい待つべきか?
数秒で十分です。1秒から始めて指数関数的バックオフを推奨し、上限を約20秒に設定してください。また、リトライ回数を制限し、上限を超えた場合は上位層に明確な失敗を返してください。
ストリーミングリクエストが切断された場合、受信済みのコンテンツは破棄すべきか?
破棄しないでください。受信済みのテキストを保持し、assistantメッセージとして追加リクエストを行うことができます。JSONなどの厳密な構造の場合、全体を再リクエストする方が安全です。
各リクエストの料金を確認するにはどうすればよいか?
レスポンス内のusageを読み取り、入力トークンに100万トークンあたり0.25ドル、出力トークンに100万トークンあたり1.00ドルを乗算することで概算できます。ストリーミング時は、usageは最後のチャンクに含まれます。
複数マシンで1つのAPI キーを共有する場合、レート制限はどのように計算されるか?
API キー単位で集計されます。すべてのマシンのリクエスト合計が毎分300回を超えないようにしてください。共有カウントまたは統一された出口キューを使用して調整する必要があります。
フォームに記入するだけでAPI キーを取得できます
アカウントを作成し、API キーをコピーし、Base URLを変更します。設定はこれだけです。