Skip to content

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 AsyncAnthropic once 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 with so 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.
Where each concern is handled in an async LLM client 4 stacked layers. Where each concern is handled in an async LLM client caller asyncio.timeout deadline; consumes the stream in async with gateway token window (95% of limit) then Semaphore(concurrency) application retries status 200 mid-stream errors; discard partial text SDK HTTP pool, SSE parsing, retries of 408/409/429/5xx before the stream The SDK's retries end where the stream begins; everything above that line is yours.

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.

The measurements behind this section A grid of 6 rows by 3 columns. The measurements behind this section behaviour measured guide unbounded gather vs Semaphore(20) 140 failed vs 0 failed of 200 fan-out semaphore only vs token window 24 failed vs 0 failed of 60 tokens per minute unclosed stream after reading 100 server generated all 2,000 cancellation error event after the stream began APIStatusError, status 200, not retried retries SDK retries at 30% overload 70% -> 97.3% success with 2 retries client CPU per streamed token 52-97 us streaming anthropic 1.11.0 and Python 3.14 against a local mock of the Messages API.

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.

How should this LLM call be configured? A decision on Who is waiting for the result with 4 outcomes. How should this LLM call be configured? Who is waiting for the result? a user, right now stream, max_retries 0-1, short deadline reserved budget share a job, minutes are fine patient retries + budget, write as you go leftover budget nobody for hours provider batch interface cheapest, slowest both at once one window and limit per key two lanes over one client Separate lanes keep a backlog from turning into a user-visible outage.

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_retries of 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.