Concurrent LLM API Calls with asyncio¶
Large language model APIs are, from asyncio's point of view, slow HTTP calls with unusual properties: responses stream for seconds or minutes, the first byte can take longer than a normal read timeout, rate limits are counted in tokens rather than requests, errors can arrive after a 200 OK, and work that nobody reads keeps running on the provider's side until the connection closes. Each property breaks a habit that works for ordinary APIs. This section measures them with the official anthropic Python SDK (1.11.0) on Python 3.14, against a local mock of the Messages API that reproduces its streaming format, its 429 and 529 responses and a configurable concurrency and token budget — so the client-side behaviour is the real SDK's, and the server's limits are known exactly. The findings: unbounded asyncio.gather over 200 calls failed 140 against a 20-request limit, where Semaphore(20) completed all 200 with zero 429s; a concurrency cap alone failed 24 of 60 calls against a token budget that a matching client-side window satisfied completely; an unclosed stream kept generating all 2,000 tokens after the client had read 100; a mid-stream error surfaced as APIStatusError with status 200 and was not retried; and streaming cost 52–97 µs of client CPU per token.
The parent section, Concurrent Execution & Worker Patterns, covers the general tools — semaphores, rate limiters, worker pools — that this section applies to one demanding kind of dependency.
Architectural principles¶
- One client, many requests. Create
AsyncAnthropiconce per process and share it; its pool and retries are per client. - Two limits, both enforced client-side. Concurrent requests with a semaphore, tokens per minute with a window or bucket whose worst case fits the provider's limit.
- Every stream has exactly one owner, which consumes it inside
async withso that every exit closes it. - Deadlines live in the caller, around the code that iterates a stream, never inside an async generator across its
yields. - Retries are layered. The SDK retries failures before a stream starts; application code retries mid-stream failures, discards partial output, and stays within a budget.
Execution model: long-lived streams on one event loop¶
A streamed LLM call is one HTTP request whose response body arrives as server-sent events over seconds or minutes. On the client it is a task suspended in await on a socket read most of the time, waking for each event: decode bytes, split lines, parse JSON, build a typed event, hand text to your code. Measured with 50 concurrent streams, that wake-up cost 52–54 µs per token with hand-written SSE parsing over httpx and 77–97 µs through the SDK's streaming helper, and event-loop lag p99 stayed between 2.4 and 12.2 ms. The arithmetic is simple: 10,000 tokens per second across all streams is half a core to a full core of event-loop time, so services that relay many streams need several processes, not just more concurrency.
The other half of the model is on the server. Generation continues until the response completes or the connection closes, and nothing else signals "stop". A client that stops reading but keeps the connection — an abandoned iterator, a background producer whose consumer left — keeps the provider generating; one that closes the connection stops it within a token or two. Timeouts are per read, not per call: the SDK's default read timeout is 600 s, which let a request with a 6-second first token complete where httpx's 5-second default failed, but a stream that keeps producing tokens is never cut off by it. Total duration needs its own deadline.
Pattern catalogue¶
Streaming with an owner and a deadline¶
Consume each stream inside messages.stream(); translate failures inside the generator; put the deadline in the caller:
async def stream_text(prompt: str):
try:
async with client.messages.stream(model=MODEL, max_tokens=1024,
messages=[{"role": "user", "content": prompt}]) as s:
async for text in s.text_stream:
yield text
except anthropic.APIStatusError as exc:
raise LLMError(exc.status_code, str(exc)) from exc
except (anthropic.APIConnectionError, httpx2.TransportError) as exc:
raise LLMUnavailable(str(exc)) from exc
async def answer(prompt: str, deadline: float = 120.0) -> str:
async with asyncio.timeout(deadline):
return "".join([t async for t in stream_text(prompt)])
A per-request timeout= that expired mid-stream raised httpx2.ReadTimeout, not an SDK exception; an asyncio.timeout placed inside the generator surfaced as a bare CancelledError for a slow consumer. See streaming LLM tokens asynchronously.
Token-based rate limiting¶
Reserve estimated input plus max_tokens before each call, against a limiter whose worst-case window spend is below the provider's limit:
await window.acquire(estimate_input_tokens(messages) + max_tokens)
A full-burst token bucket with the correct average rate still drew 60 429s; a matching sliding window drew none, and needed a 5% margin to stay at zero under heavier load. See rate limiting LLM calls by tokens per minute.
Cancellation that reaches the server¶
Every exit from the async with — break, exception, task cancellation — closes the HTTP stream, and the server stops:
async with client.messages.stream(**req) as stream:
async for text in stream.text_stream:
if enough(text):
break # server stopped at ~100 tokens of 2,000
See cancelling streaming LLM responses.
Bounded fan-out¶
A semaphore at the provider's concurrent-request limit and per-item error capture:
async with asyncio.TaskGroup() as tg:
tasks = [tg.create_task(classify_one(i, t)) for i, t in enumerate(tickets)]
See fanning out LLM calls with bounded concurrency.
Layered retries¶
SDK retries for pre-stream failures — max_retries small for interactive callers — and an application loop for mid-stream ones, with a budget during outages. See retrying LLM API calls on 429 and overload errors.
Relays that stop with their reader¶
The upstream stream opened inside the StreamingResponse generator stopped at 50 tokens when the browser left; a background producer kept generating. See relaying LLM token streams through FastAPI.
Interactive and batch traffic need different settings¶
The same gateway usually serves two kinds of caller with opposite priorities. Interactive requests — a chat reply, an autocomplete, an agent step a user is watching — care about time to first token and total latency; a failure after a few seconds with a clear message is better than a long silence. Batch work — classifying a backlog, summarizing a document set, nightly enrichment — cares about completion and cost; it can wait minutes for capacity and should never fail an item that a later retry would have saved.
The settings follow from that. Interactive calls stream, use max_retries of 0 or 1 inside a deadline of a few tens of seconds, and get a reserved share of the concurrency and token budgets so a batch job cannot starve them. Batch calls need not stream, retry more patiently with a retry budget, write results as they finish so a crash loses nothing, and take whatever budget the interactive share leaves. Two gateways over one client, or one gateway with two priority lanes, express this directly; the lane idea is the same as worker pools with priority lanes. When the provider offers an asynchronous batch interface with a lower price and a longer turnaround, it is usually the better home for work that can wait hours.
Resource boundaries¶
Five limits keep an LLM-calling service inside its budget and its provider's rules:
- Concurrent requests: a semaphore at the provider's limit per key, divided among processes. Each slot is an open HTTP stream.
- Tokens per window: reservations of input plus
max_tokens, against 95% of the provider's limit to absorb clock skew between client and server. - Streams per user: two or three concurrent answers per user, rejected with 429 beyond that, so one user cannot consume the process's share of streams, tokens and CPU.
- CPU per process: at 52–97 µs per token, the token rate a process relays determines how many processes you need.
- Retries:
max_retriesof 1 for interactive paths and around 6 for batch, plus a budget that caps retries as a fraction of first attempts during provider incidents.
Integrated production example¶
A gateway that every LLM call in a process goes through. It reserves tokens from a sliding window with a safety margin, takes a concurrency slot, consumes the stream inside a context manager, and retries both pre-stream and mid-stream failures itself — the SDK's own retries are turned off so that every attempt is counted against the token window:
import asyncio
import collections
import json
import random
import time
import anthropic
import httpx2
from anthropic import AsyncAnthropic
RETRYABLE = {200, 429, 500, 502, 503, 529} # 200: an error event after the stream began
class SlidingWindow:
def __init__(self, limit: int, window: float) -> None:
self.limit, self.window = limit, window
self.spent: collections.deque[tuple[float, int]] = collections.deque()
self._lock = asyncio.Lock()
async def acquire(self, n: int) -> None:
async with self._lock:
while True:
now = time.monotonic()
while self.spent and now - self.spent[0][0] >= self.window:
self.spent.popleft()
if sum(t for _, t in self.spent) + n <= self.limit:
self.spent.append((now, n))
return
await asyncio.sleep(self.spent[0][0] + self.window - now + 0.01)
class LLMGateway:
def __init__(self, client: AsyncAnthropic, *, tokens: int, window: float,
concurrency: int, attempts: int = 5) -> None:
self.client = client.with_options(max_retries=0) # every attempt passes the limiter
self.tokens = SlidingWindow(int(tokens * 0.95), window + 0.25) # margin for clock skew
self.slots = asyncio.Semaphore(concurrency)
self.attempts = attempts
@staticmethod
def _estimate(messages: list[dict], max_tokens: int) -> int:
return max(1, len(json.dumps(messages)) // 4) + max_tokens
async def complete(self, prompt: str, *, model: str = "claude-sonnet-5-5",
max_tokens: int = 200) -> str:
messages = [{"role": "user", "content": prompt}]
for attempt in range(self.attempts):
await self.tokens.acquire(self._estimate(messages, max_tokens))
parts: list[str] = []
try:
async with self.slots:
async with self.client.messages.stream(model=model, max_tokens=max_tokens,
messages=messages) as stream:
async for text in stream.text_stream:
parts.append(text)
return "".join(parts)
except anthropic.APIStatusError as exc:
if exc.status_code not in RETRYABLE or attempt == self.attempts - 1:
raise
except (anthropic.APIConnectionError, httpx2.TransportError):
if attempt == self.attempts - 1:
raise
await asyncio.sleep(min(8.0, 0.5 * 2 ** attempt) * random.uniform(0.5, 1.0))
raise AssertionError("unreachable")
Run with 300 concurrent calls against a mock that allowed 20 concurrent requests and 20,000 tokens per 10-second window, answered 10% of requests with 529, and broke 5% of streams with a mid-stream error: all 300 completed, with 361 HTTP attempts, zero 429s, and the 44 overloads and 17 mid-stream errors recovered by the gateway's retries. Two earlier versions of the same run were instructive. With the SDK's default max_retries=2 left on, its internal retries bypassed the token window and the server answered 25 attempts with 429. With the SDK retries off but no safety margin, 29 attempts still drew token 429s, because each reservation aged out of the client's window a moment before the server's. Partial text from a failed attempt is discarded with parts, so a retried answer is never a splice of two generations.
Diagnostic hook callout¶
Per call, record time to first token, tokens per second, total duration, attempts, and output tokens reported by the provider; per process, record in-flight streams, token-window utilisation and event-loop lag. Alert on:
- 429s above zero over a sustained window — the client's limits no longer match the provider's, or another client shares the key.
- Reported output tokens far above delivered chunks — streams are not being closed when readers stop.
- Rising time to first token with steady tokens per second — queueing at the provider; consider shedding load.
- Falling tokens per second with high event-loop lag — your own process is CPU-bound relaying tokens; add processes.
- Attempts per completed call above about 1.1 outside provider incidents — retries are hiding a configuration problem.
Failure modes¶
| Failure mode | Root cause | Detection | Fix |
|---|---|---|---|
| Most of a batch fails with 429 | Unbounded gather |
Attempts per call far above 1 | Semaphore at the provider's limit |
| 429s despite a semaphore | Token budget, not concurrency, is binding | 429 rate independent of concurrency | Token window or sized bucket |
| 429s despite a token limiter | Burst exceeds the window, SDK retries bypass it, or clock skew | Reserved vs server-side counts | Size by worst window; retries through the gateway; 5% margin |
| Tokens generated after the reader left | Unclosed stream or detached producer | Reported vs delivered tokens | async with; producer cancelled with consumer |
| Mid-stream failure not retried | SDK retries stop at the first byte | APIStatusError with status 200 |
Application retry, discard partial text |
Slow-consumer deadline raises CancelledError |
asyncio.timeout inside an async generator |
Bare CancelledError in callers |
Deadline in the caller |
| Streams time out at 5 s | httpx default timeout with slow first token | ReadTimeout before the first token |
SDK defaults or a longer read timeout |
| Latency rises with stream count | Event-loop CPU per token | Loop lag at peak streams | More processes; lighter per-event work |
Frequently Asked Questions¶
How do I make concurrent LLM API calls in Python asyncio?
Share one AsyncAnthropic client, cap concurrent calls with a semaphore at the provider's limit, reserve tokens from a limiter that fits its tokens-per-minute window, and collect results with a TaskGroup. Against a 20-request limit, Semaphore(20) completed 200 calls with no 429s where unbounded gather failed 140.
Why do I still get 429 errors with a rate limiter?
In testing, three causes: a token bucket whose burst plus refill exceeded the server's window, SDK-internal retries that bypassed the limiter, and clock skew between client and server windows. Size the limiter by its worst window, route retries through it, and keep a 5% margin.
Do I need to close LLM streams in Python?
Yes. An iterator from create(stream=True) abandoned without closing let the server generate all 2,000 tokens after the client read 100; leaving async with client.messages.stream(...) closed the stream and the server stopped at 100.
Does the Anthropic SDK retry streaming errors?
It retries errors before the stream starts. An error event after streaming began raised APIStatusError with status_code 200 and was not retried; retry it in application code and discard the partial output.
How much CPU does relaying LLM token streams take?
About 52–97 µs per token in testing, depending on whether SSE was parsed by hand or by the SDK — roughly half a core to a full core per 10,000 tokens per second.
Related¶
- Streaming LLM tokens asynchronously — streams, timeouts and the exceptions they raise.
- Rate limiting LLM calls by tokens per minute — token windows that fit the provider's.
- Cancelling streaming LLM responses — every exit closes the stream.
- Fanning out LLM calls with bounded concurrency — batches without 429 storms.
- Retrying LLM API calls on 429 and overload errors — what the SDK covers and what it does not.
- Relaying LLM token streams through FastAPI — browser-facing streams that stop with their reader.
- Rate Limiting & Throttling — the general limiter patterns.
- Concurrent Execution & Worker Patterns — the parent section.