Skip to content

Fanning Out LLM Calls with Bounded Concurrency

Batch work over an LLM — summarize 200 documents, classify 10,000 tickets, extract fields from every PDF in a folder — is a fan-out: many independent calls whose results are collected at the end. asyncio.gather over all of them is the obvious code and the wrong one, because providers cap concurrent requests and answer the excess with 429 rate_limit_error. To measure it, 200 calls of 100 output tokens each were sent through the anthropic SDK on Python 3.14 to a local mock of the Messages API that allowed 20 concurrent requests and returned 429 with retry-after: 1 beyond that. Unbounded gather with the SDK's default two retries completed only 60 calls; 140 failed, after 540 HTTP attempts of which 480 were 429s. Raising retries to 10 got all 200 through in 9.9 s — with 1,100 attempts. A Semaphore(50) failed 40. A Semaphore(20), matching the server's limit, completed all 200 in 5.6 s with 200 attempts and zero 429s; Semaphore(10) also succeeded but took 11.1 s. This guide sizes and structures LLM fan-outs.

Prerequisites

1. Measure what unbounded gather does

The tempting version launches every call at once:

async def classify_all(tickets: list[str]) -> list[str]:
    return await asyncio.gather(*(classify(t) for t in tickets))   # 200 at once

Measured against the 20-concurrent limit: 60 succeeded and 140 raised RateLimitError after exhausting their retries. Every retry wave hit the same wall — 200 callers competing for 20 slots, each told to come back in one second, all coming back together. Raising max_retries to 10 eventually pushed everything through in 9.9 s, but with 900 rejected attempts, which on a real API is quota consumed by failures and noise in the provider's view of your account. gather also has the default failure semantics of failing fast on the first exception that is not caught inside the coroutine, which for a batch job means one failed call can abandon the rest — see migrating from gather to TaskGroup.

Verify: count HTTP attempts per completed call in a batch run; anything well above 1.0 means the client is discovering the limit by hitting it.

200 calls against a 20-concurrent-request limit A grid of 5 rows by 5 columns. 200 calls against a 20-concurrent-request limit strategy completed HTTP attempts 429s time gather, no limit, retries=2 60 of 200 540 480 2.7 s gather, no limit, retries=10 200 1,100 900 9.9 s Semaphore(50) 160 320 160 4.6 s Semaphore(20), the server's limit 200 200 0 5.6 s Semaphore(10) 200 200 0 11.1 s Local mock: 100 output tokens per call at 5 ms each; 429 with retry-after 1 above 20 in flight.

2. Bound concurrency at the provider's limit

A semaphore at or slightly below the provider's concurrent-request limit lets every call through on its first attempt:

sem = asyncio.Semaphore(20)


async def classify(ticket: str) -> str:
    async with sem:
        msg = await client.messages.create(
            model="claude-haiku-4-5-20251001", max_tokens=100,
            messages=[{"role": "user", "content": f"Classify: {ticket}"}],
        )
    return msg.content[0].text

Measured: 200 of 200 in 5.6 s with exactly 200 attempts. At 10 slots the same work took 11.1 s — throughput scales with the concurrency the provider allows, so do not leave it on the table. Above the limit throughput did not improve and failures appeared: the 50-slot run was faster than the 20-slot one only because 40 of its calls failed quickly. Where the provider's limit is not published per key, find it by raising the semaphore until 429s appear and backing off 10–20%. If several processes share an API key, divide the limit between them, as with per-tenant rate limits.

Verify: at the chosen concurrency, a full batch completes with zero 429s and one attempt per call.

3. Collect results and failures without losing either

A batch should finish even when some calls fail permanently, and report which. Catch per item inside the bounded call, and keep results aligned with inputs:

from dataclasses import dataclass


@dataclass
class Result:
    index: int
    text: str | None = None
    error: str | None = None


async def classify_one(i: int, ticket: str) -> Result:
    async with sem:
        try:
            msg = await client.messages.create(model=MODEL, max_tokens=100,
                                               messages=[{"role": "user", "content": ticket}])
            return Result(i, text=msg.content[0].text)
        except anthropic.APIStatusError as exc:
            return Result(i, error=f"{exc.status_code}: {exc.message}")


async def classify_all(tickets: list[str]) -> list[Result]:
    async with asyncio.TaskGroup() as tg:
        tasks = [tg.create_task(classify_one(i, t)) for i, t in enumerate(tickets)]
    return [t.result() for t in tasks]

Catching APIStatusError per item turns a permanent failure — a 400 for a malformed prompt, a 429 that outlasted its retries — into a recorded result rather than an exception that cancels the group. Anything else, such as a bug or cancellation, still propagates and stops the batch, which is what you want. Keeping the index lets the caller retry only the failed items in a second pass.

Verify: a batch containing one invalid prompt completes, with exactly one Result carrying an error.

A bounded, failure-tolerant LLM fan-out A flow of 5 stages. A bounded, failure-tolerant LLM fan-out enumerate inputs keep the index TaskGroup task per item tasks are cheap Semaphore(N) slot N = provider limit call + SDK retries 429 rarely seen Result(index, text or error) in input order Tasks are unlimited; concurrent API calls are not.

4. Stream results to disk as they finish

For large batches, waiting for all results before writing anything wastes the work done if the process dies at item 9,000. Write each result as it completes, and make the job resumable:

async def run_batch(items: list[tuple[str, str]], out_path: str, done: set[str]) -> None:
    lock = asyncio.Lock()
    with open(out_path, "a") as out:
        async def one(item_id: str, text: str) -> None:
            if item_id in done:
                return                                   # finished in an earlier run
            async with sem:
                msg = await client.messages.create(model=MODEL, max_tokens=100,
                                                   messages=[{"role": "user", "content": text}])
            line = json.dumps({"id": item_id, "label": msg.content[0].text})
            async with lock:
                out.write(line + "\n")
                out.flush()

        async with asyncio.TaskGroup() as tg:
            for item_id, text in items:
                tg.create_task(one(item_id, text))

On restart, read the output file to rebuild done and skip completed items. Creating one task per item is fine up to tens of thousands — each costs a few kilobytes while it waits on the semaphore — but for millions of items, feed a fixed set of worker tasks from a queue instead, as in building a staged async pipeline with bounded queues. Provider batch APIs, where available, are another option for work that can wait hours rather than minutes.

Verify: killing the batch halfway and restarting it completes the remaining items without repeating finished ones.

5. Combine the concurrency cap with a token budget

Concurrency is one of two limits. A batch of long prompts can stay within 20 concurrent requests and still exceed tokens per minute. Gate each call on both, tokens first:

async def call(gateway, **kwargs):
    reserve = estimate_input_tokens(kwargs["messages"]) + kwargs["max_tokens"]
    await gateway.tokens.acquire(reserve)              # wait for token budget
    async with gateway.slots:                          # then for a concurrency slot
        return await client.messages.create(**kwargs)

Waiting for tokens before taking a slot keeps an expensive call that must wait for budget from holding a slot a cheaper call could use. The token limiter's design — and why a naive token bucket still drew 429s — is in rate limiting LLM calls by tokens per minute. Together the two limits make the client's traffic fit the provider's rules by construction, leaving retries for genuine transient errors.

Verify: a batch of maximum-size prompts runs with zero 429s, with time spent waiting on tokens visible in the limiter's metrics.

How should this LLM batch be run? A decision on How big is the batch with 4 outcomes. How should this LLM batch be run? How big is the batch? hundreds of short calls TaskGroup + Semaphore(limit) 200 in 5.6 s long prompts add a token limiter, tokens first fits TPM too thousands, hours of work write as you go, resume no lost work millions worker pool + queue, or a batch API bounded memory The semaphore turns 480 rejected attempts into zero.

Verification

An LLM fan-out is well-behaved when:

  • Concurrent calls are capped at the provider's limit (shared between processes on one key).
  • Attempts per completed call are close to 1.0, and 429s are rare.
  • Per-item failures are recorded with their index, and the batch completes.
  • Long batches write results incrementally and can resume.

Diagnostic Hook: chart in-flight calls, completed calls per second and 429s per second during a batch. A flat in-flight line at the semaphore size with steady completions is healthy; in-flight pinned while completions fall means the provider is slowing down, not rejecting — lower concurrency is unlikely to help, and the token budget may be the binding limit.

Pitfalls & edge cases

  • Unbounded gather. Measured: 140 of 200 failed, with 480 rejected attempts.
  • More retries instead of a limit. Measured: all succeeded, at 1,100 attempts for 200 calls.
  • A semaphore above the limit. Measured: 40 failures at 50 slots.
  • Exceptions that cancel the batch. Catch API errors per item; let bugs propagate.

Frequently Asked Questions

How many concurrent requests should I send to an LLM API?

The provider's concurrent-request limit for your key, divided among processes that share it. In testing against a 20-request limit, Semaphore(20) completed 200 calls with no 429s in 5.6 s; Semaphore(10) took 11.1 s and Semaphore(50) had 40 failures.

Why does asyncio.gather cause rate limit errors with LLM APIs?

It starts every call at once, so all but the allowed number are rejected, and their retries collide again. Unbounded gather failed 140 of 200 calls even with the SDK's retries.

Should I just increase max_retries for 429 errors?

It hides the problem: ten retries completed all 200 calls but needed 1,100 attempts. A semaphore at the provider's limit completed them with 200.

How do I keep a large LLM batch from losing work if it crashes?

Write each result as it completes, keyed by an item ID, and skip IDs already in the output when restarting.