Stabilność API pośredniczącego: timeouty, retry, kolejki limitów i przerwania strumienia
Po podłączeniu API pośredniczącego do produkcji, najwięcej wysiłku pochłaniają nie integracje, ale rzadkie wyjątki: timeouty w kilku zapytaniach, limit zapytań w trakcie batcha lub przerwanie strumienia. Artykuł krok po kroku omawia rozwiązania: klasyfikacja błędów, kod retry z backoffem, kolejki równoległe, kontynuacja strumienia i monitorowanie zużycia.
Kluczowe wnioski
- Najpierw klasyfikuj, potem retry: 429 i 503 wymagają backoffu, 400/401/402/403 to marnowanie czasu przy retry.
- Limit 300 zapytań na minutę należy pilnować na dwóch poziomach: ograniczania równoległych zapytań za pomocą semafora oraz kontrolowania częstotliwości; sama jedna z tych warstw nie wystarczy.
- W strumieniowaniu ustaw timeout odczytu jako „maksymalny odstęp między blokami danych”. Po przerwaniu zachowaj odebrane dane i kontynuuj.
- Zapisywanie usage dla każdego zapytania pozwala wcześnie wykryć nieprawidłowości kosztów i rozrost promptów.
Najpierw podziel błędy na trzy kategorie
Pierwszym krokiem w przypadku problemów ze stabilnością nie jest zwiększanie liczby prób ponownych, ale wiedza, co należy ponowić. Metodę postępowania można podzielić, a typowe niepowodzenia sklasyfikować w trzy kategorie:
| Kategoria | Typowe objawy | Metoda postępowania |
|---|---|---|
| Przejściowe i możliwe do odtworzenia | 429 limit zapytań, 503 upstream_busy, reset połączenia, timeout | Retry z backoffem eksponencjalnym, z limitem liczby prób |
| Błąd w żądaniu | 400 (input + max_tokens przekracza 100 000, za duże ciało zapytania), 403 content_blocked | Bez retry, popraw żądanie lub zwróć błąd użytkownikowi |
| Problem ze stanem konta | 401 nieprawidłowy klucz API, 402 brak kredytu | Bez retry, alert do właściciela, doładuj lub zmień klucz |
Tabelę tę warto umieścić w komentarzu do kodu. Typowy scenariusz awarii: gdy saldo się wyczerpie, endpoint zaczyna zwracać kod 402, a Twoja logika ponawiania nie rozróżnia kodów stanu, przez co każde zapytanie jest ponawiane pięciokrotnie, co niepotrzebnie zwiększa ruch sieciowy pięciokrotnie i zamienia logi w czerwoną masę.
Ustaw timeouty warstwowo
Sam jeden timeout całkowity nie wystarczy. Zalecamy podział na trzy warstwy:
- Timeout połączenia: Czas na nawiązanie TCP i TLS. Zazwyczaj 5–10 sekund wystarczy. Jeśli połączenie nie dojdzie do skutku, błąd powinien nastąpić szybko.
- Limit czasu na odczyt: czas oczekiwania na dane z odpowiedzi. W zapytaniach niestreamingowych czeka się na wygenerowanie całego tekstu, więc limit musi obejmować najdłuższą możliwą odpowiedź; w zapytaniu strumieniowym staje się to „maksymalną ciszą między dwoma fragmentami danych”, gdzie 10–30 sekund jest rozsądne.
- Timeout biznesowy całkowity: Maksymalny czas, jaki Twoja aplikacja może znieść. Zapewnij go przez
asyncio.wait_forlub warstwę bramki, wliczając w to wszystkie retry.
Im dłuższy output, tym dłuższy czas w trybie niestreamingowym. Jeśli ustawisz max_tokens na limit 32 000 i ustawisz tylko 30 sekund timeoutu, przy długich tekstach będziesz często dostawać timeouty. Bezpieczniej jest używać strumieniowania dla długich outputów – daje to natychmiastową informację zwrotną i jasno definiuje timeout. Parametr timeout w SDK może przyjmować liczbę lub konfigurację szczegółową – sprawdź dokumentację dla swojej wersji.
Retry z backoffem eksponencjalnym dla 429 i 503
Istota backoffu to trzy zasady: opóźnienie rośnie wykładniczo wraz liczbą niepowodzeń, należy ustawić górny limit oraz dodać losowy jitter. Bez jitteru grupa zadań, które nie powiedą się jednocześnie, ponowi próbę w tym samym momencie i ponownie obciąży właśnie odzyskany serwis. Poniższa funkcja wyłącza wbudowane ponawianie w SDK, umożliwiając nam jednolite zarządzanie retry’ami; ponawiamy tylko 429, 503 oraz błędy połączeniowe:
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 oznacza tymczasowe zajęcie serwera; zalecamy retry po kilku sekundach. Powyższe ustawienia (bazowe opóźnienie 1s, limit 20s, max 5 prób) wystarczają. Jeśli 429 występuje często, problemem nie jest retry, ale brak kontroli szybkości po stronie nadawcy – patrz następny rozdział.
Kolejka równoległych zapytań przy limicie 300 na minutę
Limit 300 zapytań na minutę na klucz to średnio 5 zapytań na sekundę. Zadania batchowe łatwo tu o błąd: wysłanie tysięcy zapytań naraz przez asyncio.gather zużywa limit w pierwszych dziesiątkach, a reszta dostaje 429.
Najlepszym rozwiązaniem jest podwójna kontrola. Semafor ogranicza liczbę „równoległych w locie” zapytań, chroniąc pamięć i połączenia lokalne. Regulator ogranicza „szybkość wysyłania”, rozkładając zapytania w czasie. Obie warstwy są konieczne: sam semafor może przekroczyć 300/min przy szybkich zapytaniach; sam regulator może gromadzić zapytania przy wolnych zapytaniach. W przykładzie ustawiamy 270 zapytań/min, zostawiając 10% zapasu na multi-instance, retry i błędy zegara:
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())
Jeśli wiele procesów lub maszyn dzieli jeden klucz, lokalny regulator nie wystarczy – limit jest sumowany dla klucza. Wtedy potrzebny jest współdzielony token bucket, np. zliczany przez Redis, albo usługa kolejki, która centralnie obsługuje wszystkie żądania batchowe.
Co zrobić po przerwaniu strumieniowania?
Połączenie strumieniowe jest bardziej wrażliwe: trwa dłużej, przechodzi przez więcej urządzeń sieciowych, a dowolny timeout bezczynności może je zerwać. Zasada: „odebrane dane to aktywo” – nie odrzucaj ich wszystkich tylko dlatego, że połączenie się urwało.
Poniższa implementacja kumuluje już otrzymane fragmenty, po przechwyceniu wyjątków połączeniowych i limitu czasu ponownie żąda istniejący tekst jako wiadomość od asystenta, aby model kontynuował pisanie. Jest to próba naprawy „best-effort”; łączenie kontynuacji nie musi być idealne, więc metoda nadaje się do długich tekstów i rozmów, ale nie do wyjść wymagających sztywnych struktur, takich jak JSON. W przypadku przerwania struktury JSON bezpieczniej jest wykonać ponowne zapytanie od początku. Zwróć uwagę, że na końcu strumieniowania automatycznie pojawia się fragment usage, który również pobieramy w kodzie:
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)
Pamiętaj, że podczas kontynuacji już wysłany tekst liczy się jako tokeny wejściowe, więc każda kontynuacja generuje dodatkowy koszt; dlatego konieczne jest ograniczenie liczby kontynuacji max_resume.
Degradacja, bezpiecznik obciążeniowy i weryfikacja przed wdrożeniem
Ponawianie rozwiązuje chwilowe awarie; jeśli awaria trwa kilka minut, dalsze ponawianie tylko gromadzi zapytania. Zalecamy dodanie mechanizmu opisu (circuit breaker) nad ponawianiem: gdy liczba kolejnych niepowodzeń przekroczy próg, np. dziesięć niepowodzeń w ciągu minuty, wstrzymaj wysyłanie zapytań na krótki czas, stosując logikę awaryjną, np. zwracając wynik z pamięci podręcznej, informując użytkownika, by spróbował później, lub odkładając zadanie w kolejce do późniejszego przetworzenia. Po zakończeniu czasu schładzania przepuść tylko kilka zapytań kontrolnych; jeśli się powiedzą, przywróć pełny ruch.
Degradację też należy zaplanować z góry. Dla czatu zwróć uprzejmą informację o zajętości; dla zadań offline zapisz błędne wpisy do tabeli błędów i zbatchuj je po zakończeniu, zamiast czekać w nieskończoność w głównym wątku. W każdym przypadku upewnij się, że błąd nie jest „połknięty” – zostaw przynajmniej jeden wiersz logu z request id.
Oto lista kontrolna przed wdrożeniem: czy wyłączono domyślne ponawianie w SDK i zostawiono tylko jedno miejsce w Twoim kodzie; czy porzucasz zapytania przy błędach 400, 401, 402, 403; czy całkowity czas oczekiwania uwzględnia wszystkie ponowienia; czy przetwarzanie wsadowe ma kontrolę limitu zapytań zamiast wysyłania równoległych zapytań naraz; czy strumieniowanie obsługuje przerwania; czy usage jest zapisywane na dysku; czy saldo ma ustawione alerty. Przejdź przez listę punkt po punkcie, a dopiero potem przeprowadzaj testy obciążeniowe, a nie na odwrót.
Monitorowanie i rozliczanie na podstawie usage
Każda odpowiedź zawiera usage, które w strumieniowaniu pojawia się w ostatnim fragmencie. Zapisanie go wraz z nazwą funkcji to najtańsza metoda monitoringu. Cena wejścia na tej stronie to 0,25 USD za milion tokenów, a wyjścia 1,00 USD za milion tokenów; koszt każdego zapytania możesz oszacować lokalnie:
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}"]
)
Po zapisaniu danych zalecamy codzienną kontrolę trzech wskaźników: mediany tokenów wejściowych w poszczególnych funkcjach, aby wykryć niepożądane rozszerzanie się promptów; udział tokenów wyjściowych, ponieważ ich cena jest cztery razy wyższa niż wejściowych, co zwykle stanowi największą część rachunku; oraz stosunek błędów 429 do 503 — wzrost tego wskaźnika oznacza konieczność dostosowania szybkości wysyłania lub strategii ponawiania. Dodaj również alert dotyczący salda, który powiadomi osobę odpowiedzialną o konieczności doładowania konta przed wystąpieniem błędu 402. Więcej szczegółów dotyczących konfiguracji znajdziesz w konfiguracji frameworka, a odpowiedzi na często zadawane pytania znajdziesz w sekcji FAQ.
Często zadawane pytania
Czy należy natychmiast ponowić żądanie po otrzymaniu błędu 429?
Nie. Najpierw zastosuj backoff z losowym jitterem i sprawdź, czy nie przekraczasz limitu 300 zapytań na minutę. Powtarzający się błąd 429 oznacza konieczność dodania regulatora (beatnika), a nie zwiększenia liczby ponowień.
Jak długo czekać na ponowienie żądania przy błędzie 503 upstream_busy?
Wystarczają kilka sekund. Zalecamy rozpoczęcie od 1 sekundy z wykładniczym backoffem, z maksymalnym limitem około 20 sekund i ograniczeniem liczby prób; po jej wyczerpaniu zwróć wyraźny błąd do warstwy wyższej.
Co zrobić z otrzymanymi danymi, jeśli połączenie strumieniowe zostanie przerwane?
Nie porzucaj ich. Możesz zachować otrzymany tekst i poprosić model o kontynuację jako wiadomość assistant; w przypadku ścisłych struktur, takich jak JSON, bezpieczniejsze jest ponowne wysłanie całego żądania.
Jak sprawdzić koszt każdego zapytania?
Odczytaj pole usage z odpowiedzi i pomnóż liczbę tokenów wejściowych przez 0,25 USD za milion tokenów oraz liczbę tokenów wyjściowych przez 1,00 USD za milion tokenów. W trybie strumieniowania dane usage są dostępne w ostatnim fragmencie odpowiedzi.
Jak obliczyć limit zapytań, gdy wiele maszyn dzieli jeden klucz API?
Statystyki są agregowane według klucza API; suma zapytań ze wszystkich maszyn nie może przekraczać 300 na minutę. Wymaga to współdzielonego licznika lub ujednoliconej kolejki wyjściowej do koordynacji.
Wypełnij formularz, aby uzyskać klucz API
Utwórz konto, skopiuj klucz i zmień Base URL. Konfiguracja jest taka prosta.