Fan Out Concurrent LLM Calls with asyncio.gather
Awaiting retrieval and LLM calls one by one wastes seconds per request. Here's how to fan them out with asyncio.gather, bound it, and handle partial failures.
A retrieval request in Archi does not touch one source. To answer an operator’s question, it searches the logbook, runs a keyword pass over JIRA tickets, and pulls the recent portal pages, then hands the merged set to a reranker. That is four calls, all over the network, none of them depending on the others. The first version I wrote awaited them one after another. The request took as long as all four added together, even though the CPU sat idle waiting on the second call while the first one had already come back.
That is the bug this post is about, and it is not a blocking-the-loop bug. It is the opposite mistake: code that is correctly async, with every call properly awaited, that still runs slow because it never overlaps the waits.
This post is for engineers building an async FastAPI backend that makes more than one independent network call per request: retrieval fan-out, a batch of embeddings, several tool calls in an agent turn. The pattern here gets those seconds back. I will show the naive version and the asyncio.gather fix. Most of the post is about the parts that bite: bounding how many calls run at once, what happens when one of them fails, and timeouts.
await does not mean concurrent
Here is the shape of the slow version: a BM25 keyword search (BM25 is a standard keyword-ranking function), a vector search, and a rerank, each an async function that awaits the network.
async def retrieve(query: str) -> list[Hit]:
bm25 = await search_bm25(query) # ~180 ms
vectors = await search_vectors(query) # ~240 ms
rerank = await rerank_hits(query, bm25 + vectors) # ~300 ms
return rerankThis code is correct. It does not block the event loop, other requests interleave fine, and the server stays responsive under load. It is just slow for this one request, because each await fully finishes before the next one starts. The await on search_bm25 says “park me until BM25 answers”, and only when BM25 answers does the code reach the await on search_vectors. The waits happen back to back, so the latencies add up: 180 plus 240 plus 300, about 720 ms, most of it spent doing nothing.
The reranker genuinely depends on the first two (it needs their hits), so it has to come last. But search_bm25 and search_vectors do not depend on each other at all, and there is no reason for the second to wait for the first. That is the waste, and it is exactly the case gather exists for.
gather starts them together
asyncio.gather takes several awaitables, schedules them all at once, and returns their results in the order you passed them in, not the order they finished. Here is the same function rewritten:
import asyncio
async def retrieve(query: str) -> list[Hit]:
# these two do not depend on each other, so run them together
bm25, vectors = await asyncio.gather(
search_bm25(query),
search_vectors(query),
)
# the reranker needs both, so it waits for the gather
return await rerank_hits(query, bm25 + vectors)Both searches are now in flight at the same time. The await on the gather returns when the slower of the two is done, so the pair costs about 240 ms instead of 420 ms. The whole request drops from roughly 720 ms to 540 ms. Fan out three or four independent sources instead of two and the gap grows, because the sequential version keeps adding while the concurrent one stays pinned to its slowest branch.
Two details are worth internalizing:
- Results come back positional.
gather(a, b)always gives[result_a, result_b], even ifbfinished first. That is why the tuple unpacking above is safe. - It only helps for work that is genuinely independent. If call two needs call one’s output, no amount of
gatherchanges that. The dependency is real, and those calls have to run in sequence.
The whole skill is spotting which calls in a request are actually independent. In a hybrid search pipeline, the lexical and vector passes are the obvious pair. In an agent turn, it is often several tool calls the model asked for in one step.
The trap: unbounded fan-out
The first time gather works, the temptation is to point it at everything. Re-embedding a thousand documents becomes a gather of a thousand coroutines; fifty queued questions become a gather of fifty LLM calls. This falls over, and it fails in a way that looks unrelated to the code you changed.
gather does not throttle. It starts every awaitable you hand it right now, all at once, so five hundred coroutines means five hundred HTTP requests opening at the same instant. That drains the connection pool. Once the pool is empty the rest queue behind it anyway, so you paid the memory for five hundred pending tasks and got none of the parallelism. Worse, if these are calls to a model provider, you just sent five hundred requests in one burst and tripped the rate limit. Now a chunk of them come back as 429s (HTTP “Too Many Requests”).
The fix is to cap how many run concurrently with an asyncio.Semaphore. A semaphore is a counter with a fixed number of permits. Each task takes a permit before doing its work and returns it afterwards, so no more than N tasks are ever inside at the same time. Wrap each call, then gather the wrappers:
import asyncio
async def gather_bounded(limit: int, *coros):
sem = asyncio.Semaphore(limit)
async def run(coro):
async with sem: # waits here if all permits are taken
return await coro
return await asyncio.gather(*(run(c) for c in coros))
# embed 500 documents, but never more than 8 requests in flight
vectors = await gather_bounded(8, *(embed(doc) for doc in documents))All five hundred tasks still get created, but the async with sem line is a gate. Only eight get through at a time; the rest suspend cheaply at that line until a permit frees up. You keep the overlap and stop it from becoming a stampede. The right value for limit is not a guess. It is whatever your provider’s rate limit and your connection pool can actually sustain, usually a small number like 5 to 20.
When one of them fails
Concurrency makes error handling less obvious, because now several things can go wrong at the same time. By default, gather is unforgiving about this. The moment one awaitable raises, that exception propagates straight up to whoever awaited the gather. You get the one error, and the results of everything else, including the calls that succeeded, are gone.
That is often the wrong behavior for retrieval fan-out. If the JIRA search times out but the logbook and the vector store both answered, I would rather rerank the two good sets than fail the whole request over one flaky source.
Pass return_exceptions=True and gather stops short-circuiting. Instead of raising, it puts each exception into the results list in that call’s slot. You get one entry per input, some values and some exceptions, and you decide what to do with each. In this version, a failed source is logged and skipped:
results = await asyncio.gather(
search_bm25(query),
search_vectors(query),
search_jira(query),
return_exceptions=True,
)
hits = []
for source, r in zip(("bm25", "vectors", "jira"), results):
if isinstance(r, Exception):
log.warning("source %s failed: %r", source, r) # degrade, don't die
else:
hits.extend(r)One thing about the default mode surprises people. When gather short-circuits on the first exception, the other tasks are not cancelled. They keep running in the background, detached, and if one of them later fails too, you get a scary “exception was never retrieved” warning from a task nobody is awaiting.
If you want cleaner semantics, where one failure cancels its siblings and the errors are collected together, use asyncio.TaskGroup (Python 3.11+). It cancels the rest of the group on the first error and raises an ExceptionGroup:
async with asyncio.TaskGroup() as tg:
t_bm25 = tg.create_task(search_bm25(query))
t_vectors = tg.create_task(search_vectors(query))
# both are awaited at the end of the block; if either raised,
# the other is cancelled and the block raises an ExceptionGroupThe rule of thumb I use:
gather(..., return_exceptions=True)when I want partial results and can degrade around a missing source.TaskGroupwhen the calls are all-or-nothing, and a failure in one means the others are wasted work worth cancelling.
Don’t forget the timeout
Overlapping the calls means the request is now only as fast as its slowest branch, which is a problem if one branch can hang. A single retrieval source that stalls for thirty seconds pins the whole gather for thirty seconds, and the concurrency you added buys nothing. Put a deadline on each call, so a slow one fails fast and becomes a handled exception instead of a hang. On Python 3.11+, the asyncio.timeout context manager reads cleanly:
async def with_deadline(coro, seconds=2.0):
async with asyncio.timeout(seconds): # asyncio.wait_for on older Pythons
return await coro
results = await asyncio.gather(
with_deadline(search_bm25(query)),
with_deadline(search_vectors(query)),
with_deadline(search_jira(query)),
return_exceptions=True,
)Now a stuck source raises TimeoutError after two seconds and lands in the results as an exception, which the loop in the previous section already knows how to skip. Bounded concurrency, per-call deadlines, and partial-failure handling are the three things that turn a gather from a demo into something you can put on a hot path.
Failure modes worth knowing
gather does nothing for CPU-bound work. It overlaps waiting, not computing. “Calls” that actually parse files or run a local embedding model on the CPU hold the GIL (Python’s global interpreter lock), so they run one at a time no matter how many you gather. That is a process pool problem, covered in the event-loop post, not a gather one.
Results are ordered, completion is not. gather hands results back in input order, so you cannot use it to stream the first answer as soon as it lands. To react to each result the moment it finishes, reach for asyncio.as_completed instead, which yields futures in completion order.
Share one client, not one per call. Creating a fresh httpx.AsyncClient inside each coroutine defeats connection pooling, because every call pays for a new TLS handshake. Build one client and pass it in, so the whole fan-out shares one pool. That client’s own pool limit is a second, quieter bound on concurrency, worth setting to match the semaphore.
A bare gather on unbounded input is a latent incident. It passes every test with ten items and falls over the day production hands it ten thousand. If the width of the fan-out comes from data rather than a fixed list, it needs the semaphore from the start, not after the first rate-limit page.
What I would do differently
Early on, I treated gather as an optimization to sprinkle on later, once something felt slow. That was backwards. The useful habit is to look at each request while writing it and ask which calls actually depend on each other. The independent ones should overlap by default, and it is easier to write the gather up front than to untangle a chain of sequential awaits after the fact. The dependency graph of a request is usually shallow: a couple of parallel fetches, then a step that needs all of them. That shape maps straight onto one bounded gather feeding one final await.
The other thing I would tell my earlier self: never ship a gather without deciding the two policies that come with it, namely how many run at once and what happens when one fails. The happy path is one line. The two seconds you save are real only if a single slow or broken source cannot take the whole request down with it. That is the difference between the fan-out in Archi and CloudCanvasAI feeling fast and it becoming the thing that pages you. That difference lives entirely in the bound and the error handling, not in the gather itself.
Diagrams by M. Hassan Ahmed, released under CC0. No external image was used for this post; the figures are original work by the author.