Handbooks / Python asyncio / Chapter 4
Async HTTP Clients & Servers
41 pages · ~70 min✓ Reviewed
Builds on Async I/O Patterns. Next up: Async Databases.
Part 1 · Async HTTP Clients and Servers in Python
Fast, Safe HTTP Calls and Servers with asyncio
Almost every real service spends its life waiting on someone else: an API that answers in 300 ms, a payment gateway, a storage bucket. With async HTTP, a coroutine sends a request and then awaits the socket, and while it waits the event loop runs other coroutines. That makes it a big win for I/O-bound fan-out, where you make many requests against high-latency endpoints. It does nothing for CPU-bound work, because the loop is still one thread.
Speed is only half the story. A client with no timeout hangs forever, a new session per request throws away pooled connections, an unbounded gather() over 10,000 URLs earns you a wave of 429 responses, and one blocking requests.get() inside an async def freezes every other request on the server. Most async HTTP bugs come from a small set of lifecycle, limit and shutdown mistakes, and this chapter walks through each of them using both aiohttp and httpx.
By the end you will be able to open one client per application and close it cleanly, size connection pools, and set per-phase and total timeouts. You will bound concurrency with a semaphore, retry transient failures with jittered backoff, and stream large bodies without loading them into memory. On the server side you will write handlers in aiohttp.web and FastAPI, keep blocking code off the loop, and shut down gracefully on SIGTERM.
You need Python 3.11 or newer, since asyncio.timeout() and TaskGroup arrive in 3.11, and a basic grasp of async def, await and asyncio.run(). To try the examples, install the libraries with pip install aiohttp httpx fastapi uvicorn. Add httpx[http2] for HTTP/2 and tenacity or aiolimiter if you want the optional retry and rate-limit helpers. Use a virtual environment so the installs stay separate from your system Python.
Part 2 · Async HTTP Fundamentals
What async HTTP actually does
An HTTP request spends almost all of its life waiting: the bytes leave your machine, travel to a server, and the answer travels back. In async HTTP your coroutine sends the request and then awaits the socket. While it waits, the event loop is free to run other coroutines, so one thread can keep many requests in flight at once.
- 1Send requestcoroutine writes to the socket
- 2await socketcoroutine pauses, loop is released
- 3Loop runs othersother coroutines make progress
- 4Reply arrivesloop wakes the coroutine
- 5Continueparse the response
This pays off when the work is I/O-bound: many requests, each with high latency. It does nothing for CPU-bound work, because the loop is still a single thread. If a coroutine spends two seconds crunching numbers, nothing else runs for those two seconds.
| Workload | Example | Does async help? |
|---|---|---|
| I/O-bound fan-out | Fetch 500 pages from slow servers | Yes: the waits overlap |
| One slow request | A single call that takes 2 s | No: there is nothing to overlap |
| CPU-bound | Hashing, parsing huge files, image resizing | No: the loop is still one thread |
Calling a coroutine function does not run it
Writing fetch(url) where fetch is an async def function only builds a coroutine object. No code inside it has run, and no request has been sent. The body starts only when something awaits the object, or when you hand it to the loop as a Task. The example below uses asyncio.sleep to stand in for network latency, so it runs anywhere.
import asyncio async def fetch(url): await asyncio.sleep(0.1) # stands in for network latency return f"body of {url}" async def main(): c = fetch("/a") # just an object, nothing has run print(type(c).__name__) print(await c) # now the body executes asyncio.run(main())
coroutine body of /a
What asyncio.run does for you
asyncio.run(main()) is the normal entry point. It creates a fresh event loop, runs main() until it finishes, cancels any tasks that are still pending, and closes the loop. The leftover-task cleanup is easy to miss, so the next example shows it happening.
import asyncio async def background(): try: await asyncio.sleep(10) except asyncio.CancelledError: print("cancelled by asyncio.run") raise async def main(): asyncio.create_task(background()) await asyncio.sleep(0) # let the task start print("main done") asyncio.run(main())
main done cancelled by asyncio.run
Calling fetch() without await sends nothing. You get a coroutine object and, if it is never awaited, a RuntimeWarning.
Sequential, gather and create_task
Putting await in a loop is the most common way to write async code that is no faster than blocking code. With for u in urls: await fetch(u), each request must finish before the next one starts, so there is no concurrency and the total time is the sum of all latencies. asyncio.gather starts every coroutine and lets the waits overlap, so the total is roughly the slowest request.
| Style | Three requests: 1 s, 2 s, 3 s | Total |
|---|---|---|
| for loop with await | One after another | 6 s (the sum) |
| asyncio.gather(...) | All waiting together | about 3 s (the slowest) |
The program below scales those latencies down to 0.1, 0.2 and 0.3 seconds and times both styles. It rounds to one decimal place so the printed numbers are stable.
import asyncio, time JOBS = [("a", 0.1), ("b", 0.2), ("c", 0.3)] async def fetch(name, latency): await asyncio.sleep(latency) return name async def sequential(): return [await fetch(n, t) for n, t in JOBS] async def concurrent(): return await asyncio.gather(*(fetch(n, t) for n, t in JOBS)) async def timed(label, make): start = time.perf_counter() result = await make() print(label, result, round(time.perf_counter() - start, 1)) async def main(): await timed("sequential", sequential) await timed("gather", concurrent) asyncio.run(main())
sequential ['a', 'b', 'c'] 0.6 gather ['a', 'b', 'c'] 0.3
create_task schedules right away
asyncio.create_task(coro) wraps a coroutine in a Task and schedules it on the loop immediately; you do not have to await it to get it started. The loop only holds a weak reference to a task, so if you drop the returned object the task can be garbage-collected while it is still running. Keep a strong reference in a set, and let a done-callback remove it when the task finishes.
import asyncio JOBS = [("a", 0.1), ("b", 0.2), ("c", 0.3)] tasks = set() async def fetch(name, latency): await asyncio.sleep(latency) return name async def main(): for name, latency in JOBS: task = asyncio.create_task(fetch(name, latency)) tasks.add(task) # strong reference task.add_done_callback(tasks.discard) # drop it when finished print("scheduled", len(tasks)) await asyncio.gather(*tasks) print("left", len(tasks)) asyncio.run(main())
scheduled 3 left 0
Dropping the result of create_task() means no reference is kept, so the task can vanish mid-flight. Store it in a set and discard it in a done-callback.
Blocking calls, aiohttp and httpx
The event loop only works if every coroutine hands control back quickly. The requests library is blocking: requests.get() holds the thread until the response arrives. Inside an async def that means the whole event loop is frozen for the length of the call, and every other coroutine stops with it. The demo below measures it with a heartbeat task that ticks every 0.05 s while a 0.12 s pause happens. time.sleep stands in for a blocking HTTP call.
import asyncio, time async def measure(label, pause): ticks = 0 async def heartbeat(): nonlocal ticks while True: await asyncio.sleep(0.05) ticks += 1 hb = asyncio.create_task(heartbeat()) await asyncio.sleep(0) # let the heartbeat start await pause() hb.cancel() print(label, ticks) async def blocking(): time.sleep(0.12) # like requests.get(): freezes the loop async def cooperative(): await asyncio.sleep(0.12) # like an async client: loop stays free async def main(): await measure("time.sleep", blocking) await measure("asyncio.sleep", cooperative) asyncio.run(main())
time.sleep 0 asyncio.sleep 2
The heartbeat could not tick at all during the blocking pause, and ticked twice during the cooperative one. Every other request your program was handling would have stalled in the same way.
Calling requests.get() or time.sleep() inside async def freezes every coroutine until it returns. Use an async client and await asyncio.sleep() instead.
Two async HTTP libraries
Two libraries do most async HTTP work in Python. aiohttp is built only for asyncio and bundles an HTTP client and a web server in one package. httpx has the same API in a sync form (Client) and an async form (AsyncClient), can speak HTTP/2 if you install httpx[http2], and runs on asyncio or trio. It is a client only.
| Feature | aiohttp | httpx |
|---|---|---|
| Async client | ClientSession | AsyncClient |
| Sync client | none | Client |
| HTTP/2 | no | yes, with httpx[http2] |
| Runtime | asyncio only | asyncio and trio |
| Web server | built in | no, client only |
Async versus threads
Threads are the older way to overlap network waits. A thread-per-request design gives each request its own OS thread with a stack that costs megabytes. With asyncio, one thread drives thousands of sockets, and each pending request is only a small task object, so a large fan-out is much cheaper. The trade-off is that the loop is shared: one blocking call stalls every request, while in a threaded design it stalls only its own thread.
| Aspect | asyncio | Threads |
|---|---|---|
| Sockets | Thousands on one thread | One thread per request |
| Cost of each | A tiny task | A stack of megabytes |
| One blocking call | Stalls every request | Stalls only that thread |
| Best for | I/O fan-out | Libraries that block |
The practical rule follows from the table: await network I/O, and never block the loop. When you are stuck with a blocking library, run that call somewhere other than the loop's thread.
Await network I/O and never block the loop. Async HTTP speeds up I/O-bound fan-out, not CPU work, and gather is what makes the waits overlap.
Part 3 · Client Lifecycle: aiohttp.ClientSession and httpx.AsyncClient
One client, opened once and closed once
An async HTTP client is more than a function that sends requests. A ClientSession in aiohttp owns a connector (the connection pool), a cookie jar and a set of default headers. An httpx.AsyncClient does the same job: a pool plus shared configuration. Because the pool is the valuable part, you create one client for the whole application and reuse it for every request, never one per request.
Anything that is the same on every call belongs on the client. Set base_url, default headers, auth and cookies once, and each call then only names the path and whatever differs.
| aiohttp | httpx | |
|---|---|---|
| Open | ClientSession() | AsyncClient() |
| Close | await s.close() | await c.aclose() |
| Auto-close | leaving async with | leaving async with |
| Body | lazy: await r.json() | eager unless you stream |
| Bad status error | ClientResponseError | HTTPStatusError |
With aiohttp, the response is also a context manager. Leaving async with s.get(url) as r: returns the connection to the pool, so the next request can reuse it. The body is read lazily, which is why you await r.json() inside the block.
import aiohttp async def load_items(url): async with aiohttp.ClientSession() as s: async with s.get(url) as r: r.raise_for_status() data = await r.json() # connection is back in the pool here # session is closed here return data
aiohttp: both blocks release something on exit
httpx reads the whole body before c.get() returns, so you call r.json() without await. The client carries the base URL and headers, so the call itself is just a path. If you want the body chunk by chunk instead, you have to ask for a stream explicitly.
import httpx async def load_items(api, key): async with httpx.AsyncClient( base_url=api, headers={'X-Key': key} ) as c: r = await c.get('/items') r.raise_for_status() return r.json() # body already read, no await
httpx: configuration lives on the client, not on each call
If you do not use async with, you own the cleanup and must call await s.close() or await c.aclose() yourself. The next example uses only the standard library to show the shape of the pattern: the exit hook runs the close call, even though the caller never writes it.
import asyncio class Client: def __init__(self): self.closed = False print('open') async def __aenter__(self): return self async def __aexit__(self, *exc): await self.aclose() async def aclose(self): self.closed = True print('closed') async def main(): async with Client() as c: print('using', c.closed) print('after', c.closed) asyncio.run(main())
open using False closed after True
What the wrong patterns cost you
A new TCP connection plus a TLS handshake costs roughly 1 to 3 extra round trips before the first byte of your request is even sent. A pooled connection pays that price once and then skips it. Creating a session for every request throws the pool away each time, so every call pays the full handshake.
| Pattern | Handshakes | Pooled |
|---|---|---|
| session per request | every call | no |
top-level httpx.get() | every call | no |
| one client per app | first call only | yes |
The top-level httpx.get() is convenient but is only a shortcut: it builds a throwaway client, sends the one request and discards the client. It is fine in a quick script, and it is wrong anywhere you make more than a handful of calls.
Neither library raises on a 4xx or 5xx response by default. A 404 or 503 comes back as an ordinary response object, and reading its body gives you an error page where you expected data. Call r.raise_for_status() before trusting the body. aiohttp then raises ClientResponseError, and httpx raises HTTPStatusError.
import aiohttp, httpx async def aio_get(s, url): try: async with s.get(url) as r: r.raise_for_status() return await r.json() except aiohttp.ClientResponseError as e: print('bad status', e.status) async def httpx_get(c): try: r = await c.get('/items') r.raise_for_status() return r.json() except httpx.HTTPStatusError as e: print('bad status', e.response.status_code)
a bad status and a network failure are different exceptions
Opening a new ClientSession or AsyncClient inside a handler or loop pays a fresh TCP and TLS handshake every time and pools nothing. Create the client once and pass it in.
A session that is never closed prints an Unclosed client session warning and leaks its sockets. Use async with, or call close() / aclose() exactly once on shutdown.
A 500 response does not raise by itself. Without raise_for_status() your code carries on and parses an error page as if it were data.
Loops, threads and long-lived services
A client binds to the event loop it was created in. Create it inside running async code, never at import time, because no loop is running then. Never share one client across event loops or threads: each loop or thread needs its own client, and a client used from the wrong loop fails in confusing ways.
One application means one client (per event loop), one open and one close. Everything in between reuses it.
In a long-lived service such as a web server, you cannot use a single async with around the whole program. Instead, open the client at startup, keep it on the application state so every handler can reach it, and close it at shutdown. In FastAPI the lifespan function does exactly this: the code before yield runs at startup, and the code after runs when the server stops.
from contextlib import asynccontextmanager import httpx from fastapi import FastAPI, Request @asynccontextmanager async def lifespan(app): app.state.http = httpx.AsyncClient(base_url='https://api.example.com') yield await app.state.http.aclose() app = FastAPI(lifespan=lifespan) @app.get('/items') async def items(request: Request): c = request.app.state.http # same client for every request r = await c.get('/items') r.raise_for_status() return r.json()
open at startup, reuse in handlers, close after yield
- 1Server startslifespan begins
- 2Client openedstored on app state
- 3Requests servedall share one pool
- 4Server stopscode after yield runs
- 5Client closedsockets released
Creating a session at module level puts it outside the running loop, so it is bound to the wrong loop or fails outright. Build it in startup code or inside an async def.
One client holds the pool, cookies and defaults. Open it inside the running loop, use async with on responses so connections return to the pool, check the status yourself, and close the client once on shutdown.
Part 4 · Connection Pooling and Limits
Why a pool exists
Opening a connection is expensive. A fresh HTTPS request needs a TCP handshake and a TLS handshake before the first byte of your request leaves, which costs one to three round trips. Keep-alive pooling avoids that cost: after a response is read, the socket stays open and goes back to a pool. The next request to the same (scheme, host, port) takes that idle socket and skips both handshakes.
- 1Ask the poolkey is (scheme, host, port)
- 2Idle socket?reuse it, no handshakes
- 3None idle, room left?open a new socket and handshake
- 4Pool full?wait in the queue for a release
- 5Response donesocket returns to the pool
A pool belongs to a client, so everything here assumes you follow the previous section and keep one long-lived ClientSession or AsyncClient. A new client per request throws the pool away each time.
The small simulation below uses a pool capped at two sockets, like limit_per_host=2. The third request finds no free socket and waits. When the first request finishes, the third one reuses its socket instead of opening a new one.
import asyncio class Pool: def __init__(self, size): self.size = size self.opened = 0 self.idle = asyncio.Queue() async def acquire(self): if self.idle.empty() and self.opened < self.size: self.opened += 1 return self.opened return await self.idle.get() def release(self, sock): self.idle.put_nowait(sock) async def request(pool, name, work): sock = await pool.acquire() print(f'{name}: socket #{sock}') await asyncio.sleep(work) pool.release(sock) print(f'{name}: done, socket #{sock} back in pool') async def main(): pool = Pool(2) await asyncio.gather( request(pool, 'req1', 0.1), request(pool, 'req2', 0.3), request(pool, 'req3', 0.1), ) asyncio.run(main())
A toy pool: two sockets, three requests
req1: socket #1 req2: socket #2 req1: done, socket #1 back in pool req3: socket #1 req3: done, socket #1 back in pool req2: done, socket #2 back in pool
Only two sockets were ever opened for three requests. req3 queued until req1 released socket #1, then ran on it with no new handshake.
The knobs in aiohttp and httpx
In aiohttp the pool lives in the TCPConnector. Its defaults are limit=100 and limit_per_host=0, and in both a value of 0 means no cap. So by default aiohttp allows up to 100 sockets in total and any number of them to a single host. Set limit_per_host when one upstream is fragile or rate sensitive: it caps simultaneous sockets to that host, so you cannot flood it even when the global limit is high.
conn = aiohttp.TCPConnector(
limit=100, # 0 = no global cap
limit_per_host=0, # 0 = no per-host cap
ttl_dns_cache=10, # seconds
force_close=False,
)
session = aiohttp.ClientSession(connector=conn)aiohttp: every value shown is the default
aiohttp also caches DNS lookups. ttl_dns_cache=10 keeps a resolved address for ten seconds, so a burst of new connections does not hit the resolver each time. A longer value saves lookups, while a shorter one picks up address changes sooner.
httpx expresses the same idea with httpx.Limits. The defaults are max_connections=100, max_keepalive_connections=20 and keepalive_expiry=5.0. Up to 100 sockets may be open at once. Only 20 idle ones are kept for reuse, and an idle socket is dropped after five seconds.
limits = httpx.Limits(
max_connections=100,
max_keepalive_connections=20,
keepalive_expiry=5.0, # seconds
)
client = httpx.AsyncClient(limits=limits)httpx: every value shown is the default
| Knob | aiohttp | httpx |
|---|---|---|
| Global cap | limit=100 | max_connections=100 |
| Per-host cap | limit_per_host (0 = none) | none, use HTTP/2 instead |
| Idle sockets kept | no separate cap | max_keepalive_connections=20 |
| Idle lifetime | server or OS decides | keepalive_expiry=5.0 |
| Pool wait fails as | connect or total timeout | PoolTimeout |
| Turn keep-alive off | force_close=True | max_keepalive_connections=0 |
| DNS cache | ttl_dns_cache=10 | none, the OS resolver |
When the pool is full, both libraries queue the request. aiohttp counts that wait inside its connect and total timeouts. httpx raises PoolTimeout if no connection frees up within the pool timeout.
When the pool is full, and HTTP/2
A full pool is not an error by itself. New requests wait for a socket to be released, and slow responses make that wait longer. The wait is easy to miss, because in aiohttp it is charged against the same clock as the request. If the wait plus the connect plus the read exceeds your total, you get a timeout even though the server never saw the request.
HTTP/2 changes the arithmetic. With httpx.AsyncClient(http2=True), many concurrent requests to the same host are multiplexed as separate streams over one connection, so the per-host socket count stays small. It needs the extra package httpx[http2], and if the server does not speak h2 the client falls back to HTTP/1.1. aiohttp has no HTTP/2 client support.
client = httpx.AsyncClient(http2=True) replies = await asyncio.gather(*(client.get(u) for u in urls)) print(replies[0].http_version) # 'HTTP/2' when negotiated
One socket per host, many streams
HTTP/2 lets many requests share a socket, but each one is still a task holding buffers. You still need a limit on in-flight tasks.
TCPConnector(force_close=True) disables keep-alive, so every request pays a fresh TCP and TLS handshake. Use it only to work around a broken server, never as a default.
Leaks, stale sockets and task limits
The most damaging pool bug is a leaked connection. A response whose body is never read or released keeps its socket checked out. A few leaks are invisible, but each one permanently removes a slot, and once the pool is empty every later request waits for a slot that will never return. From the outside it looks like the whole client suddenly hangs. The leak below reproduces that with a two-slot pool, then shows the same calls when each connection is released.
import asyncio async def call(sem, n, leak): try: await asyncio.wait_for(sem.acquire(), 0.2) except TimeoutError: print(f'call {n}: pool exhausted after 0.2s') return print(f'call {n}: got a connection') if not leak: sem.release() async def main(): for leak in (True, False): sem = asyncio.Semaphore(2) print('leaking' if leak else 'releasing') for n in (1, 2, 3): await call(sem, n, leak) asyncio.run(main())
A leaked slot is never given back
leaking call 1: got a connection call 2: got a connection call 3: pool exhausted after 0.2s releasing call 1: got a connection call 2: got a connection call 3: got a connection
The fix in real code is to scope every response with async with, so the socket goes back to the pool on every path, including early returns and exceptions. With httpx, use c.stream() in an async with for large bodies, and always finish or close the response.
# leaks: the early return never releases the socket r = await session.get(url) if r.status != 200: return None # safe: the socket returns when the block exits async with session.get(url) as r: if r.status != 200: return None body = await r.read()
Scope every response
A second failure comes from the server side. Servers and load balancers close idle connections after some time. If your pool still holds that socket, the next request on it fails with an error such as Server disconnected. The request never reached the application, but the client cannot be sure, so retry only idempotent requests (GET, PUT, DELETE) automatically. A blind retry of a POST could run the action twice.
Finally, the pool limit caps sockets, not tasks. Create 1,000 coroutines against a pool of five and all 1,000 tasks still exist, each holding memory while it waits. The pool only decides how many of them make progress at once.
import asyncio async def main(): pool = asyncio.Semaphore(5) active = peak = 0 async def req(): nonlocal active, peak async with pool: active += 1 peak = max(peak, active) await asyncio.sleep(0.001) active -= 1 tasks = [asyncio.create_task(req()) for _ in range(1000)] print('tasks created:', len(tasks)) await asyncio.gather(*tasks) print('peak in flight:', peak) asyncio.run(main())
A pool of 5 does not stop 1,000 tasks from existing
tasks created: 1000 peak in flight: 5
| Symptom | Likely cause | Fix |
|---|---|---|
| Later calls hang | Leaked responses | Wrap each response in async with |
| PoolTimeout | Pool too small for the load | Raise the limit or bound the tasks |
| Server disconnected | Stale idle socket | Retry idempotent requests only |
| Every call is slow | force_close=True | Leave keep-alive on |
| Memory climbs with queue | 10k tasks waiting on the pool | Bound tasks with a Semaphore |
Retrying a POST after a dropped socket can create the order twice. Retry GET, PUT and DELETE, and give a POST an idempotency key before you retry it.
A pool limit caps sockets, not tasks and not requests per second. Pair it with a Semaphore on in-flight work and with timeouts, which later sections cover.
Part 5 · Per-Request and Total Timeouts
Why every request needs a timeout
A request with no timeout waits for as long as the server takes to answer, and a stalled server never answers. The await just sits there. It also holds a pooled connection and a task for as long as it hangs. One dead upstream can therefore use up your whole pool and freeze every caller behind it. So set a timeout on every request.
aiohttp does ship with a default, ClientTimeout(total=300, sock_connect=30). That is a 5 minute ceiling on the whole request and 30 seconds just to open the TCP connection. For an API call a user is waiting on, that is almost always far too long, so treat it as a safety net rather than a setting.
aiohttp: ClientTimeout
aiohttp describes its limits in one ClientTimeout object. Some fields cap the whole operation and others cap a single phase. Note that connect includes any time the request spends queued waiting for a free pooled socket, whereas sock_connect measures only the TCP handshake.
| Field | What it bounds | Notes |
|---|---|---|
| total | The whole operation, from start to the last byte | The only hard cap on elapsed time |
| connect | Pool wait plus connecting | A saturated pool can trip it |
| sock_connect | The TCP connect alone | Default is 30s |
| sock_read | Longest gap between two reads | Resets with every chunk that arrives |
The session carries a default timeout, and any single request can override it. Pass a timeout to the call, as in s.get(url, timeout=aiohttp.ClientTimeout(total=5)). The per-request value replaces the session default for that call only, so a slow report endpoint can have a generous limit while the rest of the session stays strict.
httpx timeouts are per phase
httpx takes a different approach. Its default is 5 seconds on every phase: connect, read, write and pool. You can tune the phases separately. httpx.Timeout(10.0, connect=3.0) gives 3 seconds to connect and 10 seconds to each of read, write and pool. Pass it to httpx.AsyncClient(timeout=...) or to a single call.
| aiohttp | httpx | |
|---|---|---|
| Default | total=300s, sock_connect=30s | 5s on every phase |
| Model | A total plus individual phases | Phases only |
| Hard cap on elapsed time? | Yes, total | No |
| Pool wait is limited by | connect and total | pool, raising PoolTimeout |
The key point is that an httpx timeout limits each phase, not the sum of them. The read timeout restarts every time a chunk of bytes arrives. A server that sends one byte every 4 seconds never goes quiet for 5 seconds, so the read timeout never fires and the request can run for hours.
The next example imitates this with a stream that yields a byte every 0.1s. Each read is guarded by a 0.25s per-read timeout, which plays the role of the phase timeout. A separate 0.35s deadline wraps the whole loop.
import asyncio async def trickle(): for _ in range(100): await asyncio.sleep(0.1) yield b"x" async def read_all(seen): it = trickle().__aiter__() while True: try: await asyncio.wait_for(it.__anext__(), 0.25) except StopAsyncIteration: return seen.append(1) async def main(): seen = [] try: async with asyncio.timeout(0.35): await read_all(seen) except TimeoutError: print("per-read limit never fired") print("total deadline fired after", len(seen), "chunks") asyncio.run(main())
Each read is fast enough, yet the overall deadline still ends the request
per-read limit never fired
total deadline fired after 3 chunksA read timeout of 5s does not mean the request finishes within 5s. It only means no single wait for data exceeds 5s. A slow trickle can keep a request alive indefinitely.
A hard deadline on the caller
To guarantee an upper bound, put a deadline around the caller's work. On Python 3.11+ use async with asyncio.timeout(10):. On older versions use asyncio.wait_for(coro, 10). When the deadline passes, the awaiting code is cancelled and you receive a TimeoutError, whatever the phase limits say.
The deadline works well around gather(). When it expires, the cancellation reaches the gather() call, and gather() cancels every child request that is still running. Nothing is left running in the background after you have given up. In the example below, one request finishes in time and two are stuck. The deadline cancels the stuck two and reports the timeout.
import asyncio async def fetch(name, delay): try: await asyncio.sleep(delay) return f"{name} done" except asyncio.CancelledError: print(f"{name} cancelled") raise async def main(): try: async with asyncio.timeout(0.3): await asyncio.gather( fetch("a", 0.1), fetch("b", 1), fetch("c", 1)) except TimeoutError: print("deadline hit") asyncio.run(main())
The deadline cancels every unfinished child
b cancelled c cancelled deadline hit
Which exception do I catch?
The two libraries report timeouts differently. aiohttp raises asyncio.TimeoutError when total runs out and ServerTimeoutError for the socket phases. httpx raises a specific class for each phase, and all of them inherit from httpx.TimeoutException. That parent class is the one to catch when you do not care which phase failed.
| Library | Timeout that fired | Exception |
|---|---|---|
| aiohttp | total | asyncio.TimeoutError |
| aiohttp | sock_connect, sock_read | ServerTimeoutError |
| httpx | connect | ConnectTimeout |
| httpx | read | ReadTimeout |
| httpx | write | WriteTimeout |
| httpx | pool | PoolTimeout |
| Your own deadline | asyncio.timeout / wait_for | TimeoutError |
Picking the numbers
Layer the limits so each one catches a different failure. A short connect timeout fails fast on a dead host. A read timeout sized to the endpoint's SLA catches a server that accepted the request and then went silent. The total deadline on the caller catches everything else, including the trickle.
- 1Caller deadlineasyncio.timeout(10) wraps everything
- 2Connectshort: 3 to 5 seconds
- 3Writesending the request
- 4Readmatched to the endpoint SLA
Use a short connect timeout (3 to 5 seconds), a read timeout matched to the endpoint's SLA, and a total deadline on the caller.
Passing timeout=None to httpx switches off every timeout at once. A stalled server then hangs the request silently, with no error. If one call really needs more time, raise a specific phase with httpx.Timeout(10.0, read=60.0) and keep the others.
Relying on aiohttp's 300s default in production means a stuck request can hold a connection for 5 minutes. Set a real total. Also remember that phase limits alone leave the caller without a deadline.
Part 6 · Bounded Concurrency for Many Requests
Why Unbounded gather() Hurts
The usual way to fan out requests is await asyncio.gather(*(fetch(u) for u in urls)). That is fine for 20 URLs. With 10,000 URLs, gather() wraps every coroutine in a task and starts them all at once. The event loop has no notion of 'too many', so nothing slows you down until something downstream breaks.
10,000 tasks live together
every buffer and response held at once
only a few sockets exist
the rest of the tasks wait in line
each socket is a file descriptor
too many opens: OSError
the server sees a burst
answers Too Many Requests
The fix is to cap how many requests are in flight at the same moment. The simplest tool is asyncio.Semaphore(n). It hands out n permits. A coroutine that wants to run does async with sem:, takes a permit, and gives it back on exit. Once n coroutines are inside the block, the rest suspend at the async with line, before any request is sent.
The pattern is to wrap your fetch in a small bounded() function that does async with sem: return await fetch(s, url). You then hand bounded(u) to gather() instead of fetch(u). The example below uses a fake fetch that only sleeps, so it runs anywhere. It counts how many calls overlap, so you can see the cap working.
import asyncio in_flight = 0 peak = 0 async def fetch(session, url): global in_flight, peak in_flight += 1 peak = max(peak, in_flight) await asyncio.sleep(0.01) in_flight -= 1 return 'ok ' + url async def main(): sem = asyncio.Semaphore(20) async def bounded(url): async with sem: return await fetch(None, url) urls = [f'/item/{i}' for i in range(100)] results = await asyncio.gather(*(bounded(u) for u in urls)) print(len(results), 'results, peak in flight:', peak) asyncio.run(main())
100 tasks are created, but never more than 20 are inside fetch
100 results, peak in flight: 20
All 10,000 tasks still exist; they just wait at the semaphore. This is fine for thousands of URLs, but for millions or an endless stream the task objects themselves become the memory problem. Use a worker pool for that.
Worker Pool: Queue Plus N Workers
A worker pool turns the picture around. Instead of creating one task per URL and making most of them wait, you create exactly N long-lived worker tasks and feed them work through an asyncio.Queue. Each worker loops: take a URL, fetch it, mark it done. The number of tasks is N no matter how many URLs there are.
Setting maxsize on the queue is what keeps memory flat. When the queue is full, await q.put(url) suspends the producer until a worker frees a slot. So a file with a billion lines or a live stream of URLs never sits in memory all at once. Only the queue's contents and the N in-progress fetches do.
import asyncio in_flight = 0 peak = 0 async def fetch(session, url): global in_flight, peak in_flight += 1 peak = max(peak, in_flight) await asyncio.sleep(0.01) in_flight -= 1 return 'ok ' + url async def main(): q = asyncio.Queue(maxsize=10) results = [] async def worker(): while True: url = await q.get() try: results.append(await fetch(None, url)) finally: q.task_done() workers = [asyncio.create_task(worker()) for _ in range(4)] for i in range(40): await q.put(f'/item/{i}') await q.join() for w in workers: w.cancel() await asyncio.gather(*workers, return_exceptions=True) print(len(results), 'results, peak in flight:', peak) asyncio.run(main())
4 workers drain 40 URLs; the queue never holds more than 10
40 results, peak in flight: 4
Shutdown has three steps. q.join() waits until every item has been marked done. Then you cancel the idle workers, which are parked on q.get(). Finally you gather them so the cancellations are collected rather than left dangling.
If a worker takes an item and never calls task_done(), q.join() waits forever. Put the call in a finally block so a failed fetch still counts as handled.
Collecting Results and Handling Failures
Bounding the work is half the job. The other half is deciding what happens when one request fails. The four common tools behave very differently, and the differences matter most when you are running thousands of requests.
| Tool | On first failure | Other tasks | Results |
|---|---|---|---|
| TaskGroup (3.11+) | raises an ExceptionGroup | cancelled | none returned; read each task's result |
| gather() | first exception propagates | keep running, not cancelled | lost to the caller |
| gather(return_exceptions=True) | nothing raised | all finish | list mixing values and exception objects |
| as_completed() | raised when you await that item | you decide | one at a time, in completion order |
asyncio.TaskGroup gives you structured concurrency: every task started in the async with block must finish before the block exits. If one task fails, the group cancels its siblings and raises an ExceptionGroup that you catch with except*. Plain gather() is looser. The first exception reaches the caller, but the other tasks are not cancelled. They keep running in the background and burn sockets, even though nobody is waiting for their results. If you want to keep the whole batch, pass return_exceptions=True and each failure comes back as an exception object in the result list. The example shows all three side by side.
import asyncio async def job(name, delay, fail=False): await asyncio.sleep(delay) if fail: raise ValueError(name) return name async def with_group(): try: async with asyncio.TaskGroup() as tg: tg.create_task(job('a', 0.01, fail=True)) slow = tg.create_task(job('b', 0.05)) except* ValueError as eg: print('group:', [str(e) for e in eg.exceptions], 'slow cancelled:', slow.cancelled()) async def with_gather(): slow = asyncio.create_task(job('b', 0.05)) try: await asyncio.gather(job('a', 0.01, fail=True), slow) except ValueError as e: print('gather raised:', e, 'slow cancelled:', slow.cancelled()) print('slow result:', await slow) async def keep_all(): out = await asyncio.gather( job('a', 0.01, fail=True), job('b', 0.02), return_exceptions=True) print('collected:', out) async def main(): await with_group() await with_gather() await keep_all() asyncio.run(main())
same failure, three different outcomes
group: ['a'] slow cancelled: True gather raised: a slow cancelled: False slow result: b collected: [ValueError('a'), 'b']
Sometimes you do not want to wait for the whole batch. asyncio.as_completed(tasks) yields awaitables in the order they finish, not the order you submitted them. That suits a progress bar, which can tick as each result lands. It also suits early exit: stop at the first good answer and cancel the rest.
import asyncio async def job(name, delay): await asyncio.sleep(delay) return name async def main(): jobs = [job('slow', 0.03), job('fast', 0.01), job('mid', 0.02)] for i, fut in enumerate(asyncio.as_completed(jobs), 1): name = await fut print(f'{i}/3 finished: {name}') asyncio.run(main())
results arrive in completion order, not submission order
1/3 finished: fast 2/3 finished: mid 3/3 finished: slow
After a plain gather() raises, the other requests are still running and still holding pool slots. Use a TaskGroup if failures should stop the rest, or return_exceptions=True if one bad URL must not lose the batch.
Three Different Limits and How to Size n
Three separate limits get mixed up in practice. They bound different things, so setting one does not give you the others. The semaphore wraps your whole unit of work. That includes the request, the parsing and any retries inside it. The connector limit only counts open sockets. A rate limiter counts how many requests start per unit of time.
| Limit | What it bounds | Tool |
|---|---|---|
| Semaphore(n) | the whole operation: request + parse + retries | asyncio |
| limit_per_host | open sockets to one host | TCPConnector |
| token bucket | requests started per second | aiolimiter |
A rate limit is not a concurrency limit. With 20 permits and 50 ms responses you could start 400 requests a second. With 20 permits and 5 s responses you would start only 4 a second. If the upstream's quota is written as 'N requests per second', a semaphore cannot enforce it. Use a token bucket such as aiolimiter.AsyncLimiter(10, 1), which allows 10 requests per 1 second window. You can combine it with a semaphore when you need both caps.
To size n, start from what the upstream can take, not from what your machine can start. Then keep n at or below the connector's limit_per_host. If n is larger, the extra tasks pass the semaphore and then just wait for a free socket inside the pool. That wait counts against your connect and total timeouts, so requests can time out without ever being sent.
Choose n from the upstream's capacity and keep n ≤ limit_per_host. Use a Semaphore or a worker pool for concurrency, TaskGroup or return_exceptions=True for failures, and a token bucket when the quota is in requests per second.
A semaphore limits how many requests overlap, not how fast they start. Fast responses free permits quickly, so you can still exceed a per-second quota and collect 429s. Add a token bucket for that.
Part 7 · Retries with Exponential Backoff
What deserves a retry
A retry is a bet that the next attempt will go differently from the last one. That bet only pays off for transient failures, where the problem sits in the network or in an overloaded server and not in your request. Connect errors and timeouts qualify. So do the status codes 429, 502, 503 and 504, which mean the server is overloaded, throttling you, or sitting behind a gateway that hiccuped.
Errors that describe the request itself will fail identically every time. A 400 or 422 means the payload is wrong, 401 and 403 mean the credentials or permissions are wrong, and 404 means the thing is not there. Retrying these only adds load and delays the real error reaching the caller.
| Failure | Retry? | Why |
|---|---|---|
| Connect error, timeout | yes | network blip, a fresh attempt may succeed |
| 429, 503 | yes | overload or throttling; honor Retry-After |
| 502, 504 | yes | gateway hiccup or restarting upstream |
| 400, 422 | no | the payload is bad and stays bad |
| 401, 403, 404 | no | the answer will not change |
The second filter is the kind of request. Repeating a call is only safe when doing it twice has the same effect as doing it once. GET, PUT and DELETE are idempotent by definition, so they are safe to retry. POST is not: if the server created the order but the response was lost, a blind retry creates a second order. A POST may be retried only when you send an Idempotency-Key header, a unique value the server uses to recognise a repeat and return the first result.
Retrying a POST without an idempotency key. When the response is lost but the work was done, the retry charges the card or creates the record twice.
Backoff, jitter and Retry-After
Retrying immediately usually hits the same overloaded server a few milliseconds later. Exponential backoff waits longer after each failure: the delay ceiling is min(cap, base * 2**attempt). With a base of 0.5 seconds that gives 0.5, 1, 2, 4 and so on, until the cap (30 seconds here) stops the growth.
If a thousand clients fail together and all wait exactly 2 seconds, they all come back together too, which is a thundering herd. Full jitter fixes this by sleeping a random time between zero and the ceiling, await asyncio.sleep(random.uniform(0, delay)), so the retries spread out across the whole window.
import random def ceiling(attempt, base=0.5, cap=30): return min(cap, base * 2**attempt) def backoff(attempt): return random.uniform(0, ceiling(attempt)) for n in (0, 1, 2, 3, 6): print(n, ceiling(n)) random.seed(1) print(all(0 <= backoff(n) <= ceiling(n) for n in range(10)))
The ceiling doubles until the cap; the actual sleep is a random point under it.
0 0.5 1 1.0 2 2.0 3 4.0 6 30 True
| Attempt | base * 2**n | Ceiling | Actual sleep |
|---|---|---|---|
| 0 | 0.5s | 0.5s | 0 to 0.5s |
| 1 | 1s | 1s | 0 to 1s |
| 3 | 4s | 4s | 0 to 4s |
| 6 | 32s | 30s (cap) | 0 to 30s |
Sometimes the server tells you exactly how long to wait. On a 429 or 503 the Retry-After header, usually a number of seconds, beats your own computed delay. Cap it so a bad value cannot stall you for hours, and fall back to jittered backoff when the header is missing.
Put these together and you get the manual loop: try up to five times, catch only the transient exception, pick the wait, sleep, and raise at the end. The simulated fetch below fails with a 503 carrying Retry-After, then with a plain 502, then succeeds.
import asyncio, random class Transient(Exception): def __init__(self, status, retry_after=None): super().__init__(status) self.status = status self.retry_after = retry_after def backoff(attempt, base=0.1, cap=30): return random.uniform(0, min(cap, base * 2**attempt)) calls = 0 async def fetch(): global calls calls += 1 if calls == 1: raise Transient(503, retry_after=0.1) if calls == 2: raise Transient(502) return 'payload' async def get_with_retry(): for attempt in range(5): try: return await fetch() except Transient as exc: if exc.retry_after is not None: wait = min(exc.retry_after, 30) source = 'Retry-After' else: wait = backoff(attempt) source = 'backoff' print(f'attempt {attempt}: {exc.status}, using {source}') await asyncio.sleep(wait) raise RuntimeError('gave up') print(asyncio.run(get_with_retry()))
Retry-After wins when present; otherwise jittered backoff.
attempt 0: 503, using Retry-After attempt 1: 502, using backoff payload
With backoff defined, the safe POST retry is short. The key is generated once, before the loop, and sent on every attempt. The simulated server records the charge under that key, then loses the first response. The retry gets the stored result instead of charging again.
import asyncio, random, uuid def backoff(attempt, base=0.05, cap=30): return random.uniform(0, min(cap, base * 2**attempt)) charged = {} responses_to_lose = 1 async def post(key, amount): global responses_to_lose if key not in charged: charged[key] = amount if responses_to_lose: responses_to_lose -= 1 raise TimeoutError('response lost') return charged[key] async def pay(amount): key = str(uuid.uuid4()) for attempt in range(3): try: return await post(key, amount) except TimeoutError: await asyncio.sleep(backoff(attempt)) raise RuntimeError('gave up') print(asyncio.run(pay(25)), len(charged))
The same key on every attempt means one charge, not two.
25 1
Sleeping the same fixed delay in every client, or skipping the cap. Without jitter the herd returns in lockstep; without a cap, attempt 10 would sleep for minutes.
Library support for retries
You rarely have to write every line yourself, but you do need to know what each library actually retries. The answer differs a lot, and several of them leave status codes alone.
| Library | Built-in retry | What it covers |
|---|---|---|
| httpx | AsyncHTTPTransport(retries=3) | only ConnectError and ConnectTimeout, never status codes |
| httpx, statuses | none | write the loop shown earlier |
| aiohttp | none | write the loop yourself, or add aiohttp-retry |
| tenacity | @retry decorator | whatever exceptions you choose |
For httpx, pass httpx.AsyncHTTPTransport(retries=3) to the client's transport argument. It is cheap protection against connection failures, because a connection that never opened cannot have done any work. It will not retry a 503 or a read timeout, so you still need your own loop for those.
tenacity decorates an async def directly, with no wrapper code. A typical setup is @retry(stop=stop_after_attempt(5), wait=wait_random_exponential(max=10), retry=retry_if_exception_type(httpx.TransportError)). The stop argument caps attempts, wait_random_exponential gives exponential backoff with full jitter, and retry_if_exception_type limits retries to the exceptions you name.
aiohttp has no retry of its own. Either write the loop, or wrap the session with aiohttp-retry: build a RetryClient with ExponentialRetry(attempts=4) and call rc.get(url) as you would on a session.
Whichever tool you use, decide the retryable exceptions and statuses first. A decorator that retries every Exception also retries your own bugs.
Slots, deadlines and cancellation
Retries interact with the bounded concurrency from the previous section. If a task sleeps for backoff while still inside async with sem, it holds a concurrency slot and does nothing with it, and other work waits behind it. The fix is to release the semaphore after each attempt and sleep outside it.
Retries also need a ceiling in time, not only in count: five attempts that each sleep up to 30 seconds can run for minutes. Wrap the whole retry loop in asyncio.timeout() so the caller has a hard deadline. In the example below the semaphore has one slot. Task a fails twice and backs off, and task b uses the freed slot and finishes first.
import asyncio, random class Transient(Exception): pass def backoff(attempt, base=0.1, cap=30): return random.uniform(0, min(cap, base * 2**attempt)) sem = asyncio.Semaphore(1) tries = {'a': 0, 'b': 0} async def try_once(name): async with sem: await asyncio.sleep(0.01) tries[name] += 1 if name == 'a' and tries[name] < 3: raise Transient() return name async def get(name): async with asyncio.timeout(2): for attempt in range(5): try: return await try_once(name) except Transient: pass await asyncio.sleep(backoff(attempt)) raise RuntimeError('gave up') async def main(): tasks = [asyncio.create_task(get(n)) for n in ('a', 'b')] for fut in asyncio.as_completed(tasks): print(await fut, 'done') print(tries) asyncio.run(main())
The sleep sits after the semaphore block, so a backing-off task holds no slot.
b done
a done
{'a': 3, 'b': 1}- 1Acquire slotasync with sem
- 2Send requestinside the slot
- 3Release slotleave the with block
- 4Sleep with jitteroutside the semaphore
- 5Check deadlineasyncio.timeout around it all
The last rule is about cancellation. When you shut down or a timeout fires, asyncio cancels a task by raising asyncio.CancelledError inside it, often right in the middle of your backoff sleep. If a retry except swallows that exception, the task carries on retrying as if nothing happened. In the demo, swallow uses except BaseException and polite re-raises.
import asyncio async def swallow(): for attempt in range(3): try: await asyncio.sleep(0.1) except BaseException: print('swallowed', attempt) return 'finished anyway' async def polite(): for attempt in range(3): try: await asyncio.sleep(0.1) except asyncio.CancelledError: raise except Exception: pass return 'finished' async def try_cancel(fn): t = asyncio.create_task(fn()) await asyncio.sleep(0.05) t.cancel() try: print(await t) except asyncio.CancelledError: print('cancelled') async def main(): await try_cancel(swallow) await try_cancel(polite) asyncio.run(main())
The cancel request is honored only when CancelledError is allowed to propagate.
swallowed 0
finished anyway
cancelledA bare except: or except BaseException around the request. It catches CancelledError, so shutdown hangs and asyncio.timeout() never takes effect. Catch the specific transient exceptions, and re-raise CancelledError first if you must catch broadly.
Retry only transient failures on idempotent requests. Use capped exponential backoff with full jitter, let Retry-After override it, sleep outside the semaphore, put a deadline around the whole loop, and never swallow CancelledError.
Part 8 · Streaming Responses
Why stream, and how each library does it
A normal request reads the whole response body into memory before handing it to you. That is fine for a small JSON payload. It is a problem for a 2 GB download, or for a feed that never ends. Streaming reads the body a chunk at a time, so memory use stays flat however large the body is. Three cases call for it: large file downloads, SSE (server-sent events) and NDJSON (one JSON document per line) feeds.
In aiohttp the body is already lazy. r.content is a StreamReader that pulls bytes off the socket as you ask for them. You choose how to read it: iter_chunked(n) yields fixed-size pieces, iter_any() yields whatever has arrived so far, and readline() returns one line at a time.
| aiohttp reader | What you get | Best for |
|---|---|---|
iter_chunked(64 * 1024) | Fixed-size bytes, except the last | Writing files |
iter_any() | Whatever has arrived, any size | Low latency, since you never wait to fill a buffer |
readline() | One line, ending at the newline | Line protocols such as NDJSON |
async def download(s, url, out): async with s.get(url) as r: r.raise_for_status() async for chunk in r.content.iter_chunked(64 * 1024): out(chunk)
aiohttp: the body stays in the socket, not in RAM
httpx is the opposite. A plain c.get() always buffers the full body before it returns, so using it on a 2 GB file loads all 2 GB. To stream you must use c.stream(), an async context manager. Inside it you pick an iterator that matches the data you want.
| httpx iterator | Yields | Notes |
|---|---|---|
aiter_bytes() | Decoded bytes | Gzip is already undone |
aiter_text() | str chunks | Decoded with the response encoding |
aiter_lines() | One text line at a time | Suits NDJSON and SSE |
aiter_raw() | Raw bytes as sent | Still compressed if the server gzipped |
async def download(c, url, out): async with c.stream('GET', url) as r: r.raise_for_status() async for chunk in r.aiter_bytes(): out(chunk)
httpx: stream() instead of get()
Calling c.get() on a huge file and expecting it to stream. httpx buffers the entire body first, so memory spikes before you see the first byte.
Rules inside a stream
Streaming gives you a response whose body you have not read yet, and that changes a few habits. First, check the status before you iterate. If the server answered with a 500 and you start writing chunks, the error page ends up on disk as if it were your data. Call r.raise_for_status() right after the stream opens, before the loop.
Sometimes you want the error body, for logging. Inside a stream() block you can call await r.aread() to load the body on demand. Do it only on the error path, so the happy path stays streamed.
async def open_checked(c, url): async with c.stream('GET', url) as r: if r.status_code >= 400: body = await r.aread() raise RuntimeError(f'{r.status_code}: {body[:200]!r}') async for chunk in r.aiter_bytes(): yield chunk
Read the body only for error responses
Second, always leave the async with. The connection goes back to the pool only when the context exits. If you break out of the loop and keep the response object around without exiting its context, the connection stays pinned. Do that enough times and the pool drains and later requests hang. The async with form handles early break and exceptions for you, so never open a stream any other way.
Skipping raise_for_status(). The 500 HTML page is saved as your file and nobody notices until someone opens it.
Leaving a stream half-read without exiting its async with. The connection leaks and the pool slowly fills.
Files, feeds, uploads and deadlines
Writing a download to disk has a trap. A plain file.write() is blocking I/O, and calling it inside the loop freezes the event loop for every other coroutine. There are two fixes. Write through aiofiles, which pushes each write to a thread for you. Or collect chunks into a buffer and hand the batch to asyncio.to_thread, which costs fewer thread hops.
import aiofiles async def save(c, url, path): async with aiofiles.open(path, 'wb') as f: async with c.stream('GET', url) as r: r.raise_for_status() async for chunk in r.aiter_bytes(): await f.write(chunk)
aiofiles: each write is awaited, so the loop stays free
import asyncio async def save_batched(c, url, path): buf = [] size = 0 with open(path, 'wb') as f: async with c.stream('GET', url) as r: r.raise_for_status() async for chunk in r.aiter_bytes(): buf.append(chunk) size += len(chunk) if size >= 1024 * 1024: await asyncio.to_thread(f.write, b''.join(buf)) buf, size = [], 0 if buf: await asyncio.to_thread(f.write, b''.join(buf))
Batch about 1 MB, then write it in a worker thread
Timeouts behave differently on a stream than you might expect. The read timeout applies to each chunk, not to the whole download. A server that sends one byte every few seconds never trips it, so a slow stream can run for hours. If you need a hard limit, add a total deadline around the whole block with asyncio.timeout.
import asyncio async def save_with_deadline(c, url, path): async with asyncio.timeout(60): await save(c, url, path)
Total cap of 60 seconds, whatever the per-chunk timeout says
SSE and NDJSON feeds are line oriented, so aiter_lines() fits. Each non-empty line is one event. Blank lines are keep-alives, so skip them. For SSE, each payload line starts with data: , which you strip before parsing.
import json async def read_feed(c, url): async with c.stream('GET', url) as r: r.raise_for_status() async for line in r.aiter_lines(): if not line: continue if line.startswith('data: '): line = line[len('data: '):] event = json.loads(line) print(event)
NDJSON and SSE, one event per line
Uploads can stream too, so a large file never has to sit in memory. Pass an async generator as data= in aiohttp, or an async iterator as content= in httpx. The client pulls the next piece only when the socket can take it.
| aiohttp | httpx | |
|---|---|---|
| Streaming upload argument | data=<async generator> | content=<async iterator> |
| Streaming download | r.content.iter_chunked() | c.stream() with aiter_bytes() |
async def pieces(path): with open(path, 'rb') as f: while chunk := f.read(64 * 1024): yield chunk async def upload_aiohttp(s, url, path): async with s.post(url, data=pieces(path)) as r: r.raise_for_status() async def upload_httpx(c, url, path): r = await c.post(url, content=pieces(path)) r.raise_for_status()
Same generator feeds either client
Trusting the read timeout as a total limit. It resets with every chunk, so a trickling download can hang for hours unless you wrap it in asyncio.timeout.
Stream with r.content in aiohttp and c.stream() in httpx. Check the status first, write to disk without blocking, put a total deadline around it, and always exit the async with.
Part 9 · Async Web Servers: aiohttp.web and FastAPI/Starlette
aiohttp.web: handlers and a shared client
An async HTTP client is only half of the story, because the same event loop can also run the server. aiohttp.web is the smallest way to do that. You write a coroutine that takes a request and returns a response object, register it on the router under a path, and hand the application to web.run_app, which starts the loop and serves until it is stopped.
Every handler must be an async def. The framework awaits it on the event loop, so a plain def handler fails. Input comes from the request: path parameters such as {id} are in request.match_info, and the body is read with await request.json(). The body read is network I/O, which is why it needs await.
from aiohttp import web async def handle(request): return web.json_response({"ok": True}) async def item(request): item_id = request.match_info["id"] # from /items/{id} body = await request.json() # I/O, so await it return web.json_response({"id": item_id, "got": body}) app = web.Application() app.router.add_get("/x", handle) app.router.add_post("/items/{id}", item) if __name__ == "__main__": web.run_app(app)
Fragment: needs aiohttp installed and starts a server
A server that calls other services needs one long-lived ClientSession, as in the client sections. aiohttp gives you cleanup_ctx for this. It takes an async generator: the code before yield runs at startup and creates the session, and the code after yield runs at shutdown and closes it. Store the session on the app under a typed web.AppKey rather than a bare string, so editors and type checkers know what it holds.
import aiohttp from aiohttp import web CLIENT = web.AppKey("client", aiohttp.ClientSession) async def client_ctx(app): app[CLIENT] = aiohttp.ClientSession() # startup yield await app[CLIENT].close() # shutdown async def upstream(request): s = request.app[CLIENT] # same session every time async with s.get("https://example.com/api") as r: return web.json_response(await r.json()) app = web.Application() app.cleanup_ctx.append(client_ctx) app.router.add_get("/up", upstream)
Fragment: shows the wiring, run it with web.run_app(app)
Writing a plain def handler in aiohttp. Handlers must be async def. Forgetting the await app[CLIENT].close() after yield leaks sockets at shutdown.
ASGI servers: what runs where
FastAPI and Starlette do not serve sockets themselves. They are ASGI applications, and an ASGI server such as uvicorn accepts connections and calls them. Each worker process runs exactly one event loop, and every request handled by that worker shares it. Which thread your endpoint code runs on depends on how you declare it.
An async def endpoint runs directly on the event loop, so while it executes nothing else in that worker can run. It may therefore only await non-blocking I/O such as httpx, asyncpg or asyncio.sleep. A plain def endpoint is handed to the anyio threadpool, which has 40 threads by default, so a blocking call such as requests.get stalls only its own thread. Dependencies follow the same rule: a def dependency also runs in the threadpool, and an async def dependency runs on the loop.
| Endpoint | Calls | Result |
|---|---|---|
| async def | await httpx | best |
| async def | requests / time.sleep | blocks EVERY request |
| def | requests | OK up to 40 concurrent |
| def dependency | blocking I/O | runs in the threadpool too |
The effect is easy to see without any web framework. In the program below a ticker coroutine stands in for all the other requests in the worker. The first handler blocks inside the coroutine with time.sleep, and the ticker freezes until it finishes. The second handler does the same sleep in a worker thread, the way a def endpoint would, and the ticker keeps going.
import asyncio, time async def ticker(log): for i in range(6): log.append(f"tick {i}") await asyncio.sleep(0.05) async def handler_async(log): log.append("async start") time.sleep(0.3) # blocks the whole loop log.append("async end") def handler_def(log): log.append("def start") time.sleep(0.3) # blocks only a worker thread log.append("def end") async def run(kind): log = [] t = asyncio.create_task(ticker(log)) await asyncio.sleep(0.075) # let two ticks happen first if kind == "async": await handler_async(log) else: await asyncio.to_thread(handler_def, log) await t print(kind + ":", ", ".join(log)) async def main(): await run("async") await run("def") asyncio.run(main())
async: tick 0, tick 1, async start, async end, tick 2, tick 3, tick 4, tick 5 def: tick 0, tick 1, def start, tick 2, tick 3, tick 4, tick 5, def end
Calling requests.get() or time.sleep() inside an async def endpoint. It freezes the loop for every client of that worker. Use await asyncio.sleep() and an async client, or declare the endpoint with plain def.
FastAPI: lifespan client, reuse and streaming proxy
FastAPI's counterpart to cleanup_ctx is the lifespan function, an async context manager. Code before yield runs once at startup and creates the httpx.AsyncClient. Code after yield runs once at shutdown and closes it. The client is stored on app.state, so every request in that worker can reach it.
from contextlib import asynccontextmanager import httpx from fastapi import FastAPI, Request U = "https://example.com/api" @asynccontextmanager async def lifespan(app): app.state.client = httpx.AsyncClient() yield await app.state.client.aclose() app = FastAPI(lifespan=lifespan) @app.get("/a") async def a(request: Request): c = request.app.state.client # reuse, never create here r = await c.get(U) return r.json()
Fragment: needs fastapi and httpx installed
Reuse that client in every request. An AsyncClient() created inside the endpoint has an empty connection pool each time, so every call pays a fresh TCP and TLS handshake and the pooling you set up is thrown away. The endpoint above only reads the client from request.app.state, and the loop stays free while it waits for the upstream.
The same client can also proxy a large upstream body without holding it in memory. StreamingResponse accepts an async generator. Inside it you open c.stream(...) and yield each chunk as it arrives, so bytes flow from the upstream straight to the caller with no buffering.
from fastapi.responses import StreamingResponse async def up(c): async with c.stream("GET", U) as r: async for b in r.aiter_bytes(): yield b # forwarded as it arrives @app.get("/proxy") async def proxy(request: Request): c = request.app.state.client return StreamingResponse(up(c))
Fragment: continues the app defined above
Creating an AsyncClient() per request loses pooling and adds handshakes. Forgetting aclose() after yield leaks sockets. Calling c.get() on a huge upstream body buffers all of it, so use c.stream() for proxies.
Using every core: workers
One worker means one event loop, which means one CPU core for your Python code. To use more cores you run more worker processes. With uvicorn that is a flag. With gunicorn you keep its process management and tell it to use uvicorn's worker class.
uvicorn app:app --workers 4 gunicorn app:app -w 4 -k uvicorn.workers.UvicornWorker
Workers share nothing. Each one starts its own event loop and runs its own lifespan, so each creates its own httpx.AsyncClient with its own pool. That is the right design, because a client must stay on the loop that created it. It also means pool limits multiply: 4 workers with max_connections=100 can open up to 400 sockets to one upstream, so split the budget across workers.
| Setup | Event loops | HTTP clients | Cores used |
|---|---|---|---|
| uvicorn app:app | 1 | 1 | one |
| uvicorn app:app --workers 4 | 4 | 4, one per worker | up to four |
| gunicorn -k uvicorn.workers.UvicornWorker -w 4 | 4 | 4, one per worker | up to four |
Only awaits on non-blocking I/O: async def. Blocking library: plain def. One lifespan client per worker, reused by every request, and one worker per core you want to use.
Part 10 · Blocking Code in Handlers and Graceful Shutdown
When a Handler Blocks the Loop
An async server runs every async def handler on one event loop thread. If a single handler makes a blocking call such as requests.get, time.sleep or a synchronous database driver, the loop cannot switch to any other coroutine until that call returns. The symptom is latency spikes across every endpoint at once, including cheap ones like a health check. Nothing is wrong with the endpoint that looks slow in your metrics; it is only waiting behind the one that blocked.
To find the culprit, run the app with the environment variable PYTHONASYNCIODEBUG=1. asyncio then logs every callback or task step that held the loop for more than 100 ms, with the name of the offending task. Debug mode adds overhead, so use it in development or staging only.
The fix is to move the blocking call off the loop thread. The simplest tool is await asyncio.to_thread(func, *args) (Python 3.9+). It runs func in the loop's default thread executor and hands back the result, while the loop keeps serving other coroutines. The example below counts how many times a background ticker runs during 0.2 s of blocking work, first called directly and then through to_thread.
import asyncio, time async def ticker(ticks): while True: await asyncio.sleep(0.02) ticks.append(1) async def count_ticks(work): ticks = [] t = asyncio.create_task(ticker(ticks)) await asyncio.sleep(0) await work() t.cancel() return len(ticks) async def blocking(): time.sleep(0.2) async def offloaded(): await asyncio.to_thread(time.sleep, 0.2) async def main(): print('blocking ticks:', await count_ticks(blocking)) print('offloaded ticks >= 5:', await count_ticks(offloaded) >= 5) asyncio.run(main())
The ticker is starved while the loop is blocked
blocking ticks: 0 offloaded ticks >= 5: True
Putting requests.get() or time.sleep() inside async def and wondering why unrelated endpoints get slow. Use an async client and await asyncio.sleep(), or offload the blocking call.
Choosing Threads, Your Own Pool, or Processes
to_thread uses a shared default executor, so you do not control how many threads it may start. When you want a hard cap, for example to protect a legacy database that tolerates only a few connections, create your own ThreadPoolExecutor(max_workers=...) and submit work with loop.run_in_executor(pool, func, *args). Jobs beyond the cap wait in the pool's queue. The demo submits six jobs to a pool of two and records the highest number of jobs running at the same moment.
import asyncio, threading, time from concurrent.futures import ThreadPoolExecutor running = 0 peak = 0 lock = threading.Lock() def work(i): global running, peak with lock: running += 1 peak = max(peak, running) time.sleep(0.05) with lock: running -= 1 return i async def main(): loop = asyncio.get_running_loop() with ThreadPoolExecutor(max_workers=2) as pool: out = await asyncio.gather( *(loop.run_in_executor(pool, work, i) for i in range(6))) print('results:', out) print('peak threads:', peak) asyncio.run(main())
results: [0, 1, 2, 3, 4, 5] peak threads: 2
Threads help when the blocking call is waiting on I/O, because a waiting thread releases the GIL. They do not help with CPU-bound work. The GIL lets only one thread execute Python bytecode at a time, so a hashing or parsing loop on threads still competes for one core. Use a ProcessPoolExecutor, which runs the function in separate worker processes with their own interpreters. The function and its arguments must be picklable, and on Windows or macOS the entry point needs the __main__ guard.
import asyncio from concurrent.futures import ProcessPoolExecutor def squares(n): return sum(i * i for i in range(n)) async def main(): loop = asyncio.get_running_loop() with ProcessPoolExecutor(max_workers=2) as pool: out = await asyncio.gather( loop.run_in_executor(pool, squares, 1000), loop.run_in_executor(pool, squares, 10)) print(out) if __name__ == '__main__': asyncio.run(main())
[332833500, 285]
Starlette Helpers and Threads You Cannot Cancel
Inside Starlette and FastAPI you can use the framework's own helpers, which run on the anyio threadpool that plain def endpoints already use (40 threads by default). The two spellings are equivalent for this purpose.
| Tool | Runs in | Use for |
|---|---|---|
asyncio.to_thread(f, *a) | default executor | blocking I/O, any asyncio code (3.9+) |
loop.run_in_executor(pool, f, *a) | your own ThreadPoolExecutor | capping max_workers |
ProcessPoolExecutor | worker processes | CPU work, because of the GIL |
await run_in_threadpool(f) or anyio.to_thread.run_sync(f) | anyio threadpool | Starlette and FastAPI handlers |
from starlette.concurrency import run_in_threadpool @app.get('/report') async def report(): rows = await run_in_threadpool(load_report, 'q3') return {'rows': len(rows)}
load_report is a blocking function; the loop stays free while it runs
There is one trap with all thread offloading: a thread cannot be cancelled. If you wrap to_thread in asyncio.wait_for or asyncio.timeout, the timeout abandons the await and raises in your coroutine, but the worker thread keeps running to the end of its function. It still holds its slot in the pool and may still write to a database or file. The demo shows the await giving up while the thread finishes later.
import asyncio, threading, time done = threading.Event() def slow(): time.sleep(0.2) done.set() async def main(): try: await asyncio.wait_for(asyncio.to_thread(slow), 0.05) except TimeoutError: print('await abandoned') print('thread finished yet:', done.is_set()) await asyncio.sleep(0.3) print('thread finished later:', done.is_set()) asyncio.run(main())
await abandoned thread finished yet: False thread finished later: True
Pass a threading.Event into the worker and check event.is_set() inside its loop. Set the event when you time out or shut down, and the thread can return early on its own.
Using threads for CPU-heavy work and expecting a speedup. The GIL means threads do not run Python bytecode in parallel; reach for ProcessPoolExecutor.
Graceful Shutdown: Order, Signals and Deadlines
When a deploy replaces your server, the old process receives SIGTERM. A graceful shutdown uses that moment to finish work already accepted instead of dropping it. The order matters: you must stop new work arriving before you drain, and you must keep the HTTP clients open until nothing can still use them.
- 1Stop acceptingclose the listener
- 2Finish in-flightdrain requests
- 3Cancel background taskscancel + gather
- 4Close HTTP clientssession / client
- 5Exitprocess ends
Your server runtime already implements the first two steps; you supply the last three through hooks.
| Runtime | On SIGTERM | Knob |
|---|---|---|
| uvicorn | stops accepting, drains in-flight requests | --timeout-graceful-shutdown N caps the wait |
| aiohttp | drains in-flight requests | web.run_app(shutdown_timeout=60), 60 s is the default |
| aiohttp hooks | close your resources | on_shutdown and on_cleanup signals, or cleanup_ctx |
| FastAPI | uvicorn drains, then lifespan resumes | code after yield in the lifespan function |
| Kubernetes | SIGKILL after the grace period | terminationGracePeriodSeconds, 30 s by default |
Closing the ClientSession or httpx.AsyncClient in cleanup closes its pooled sockets cleanly and avoids Unclosed client session warnings and leaked connections. In aiohttp the natural place is a cleanup_ctx generator: code before yield runs at startup, code after it runs at shutdown. In FastAPI it is the part of the lifespan after yield, where you call await client.aclose().
from aiohttp import ClientSession, web HTTP = web.AppKey('http', ClientSession) async def http_ctx(app): app[HTTP] = ClientSession() yield await app[HTTP].close() app = web.Application() app.cleanup_ctx.append(http_ctx) web.run_app(app, shutdown_timeout=25)
The drain window is set below the Kubernetes grace period
The drain window has to fit inside the platform's deadline. Kubernetes sends SIGTERM and starts a clock; when terminationGracePeriodSeconds runs out it sends SIGKILL, which ends the process mid-request. So keep your total shutdown budget under that value.
| Phase | Budget |
|---|---|
| Drain in-flight requests | up to 20 s |
| Cancel background tasks | up to 3 s |
| Close HTTP clients | up to 2 s |
| Total | 25 s, below the 30 s grace period |
Cancelling Background Tasks and Shutdown Pitfalls
Fire-and-forget tasks keep running after the last request is drained, so they need an explicit stop. Keep every background task in a set (which also stops it being garbage collected mid-flight), and remove it with a done callback. At shutdown, call task.cancel() on each one, then await asyncio.gather(*tasks, return_exceptions=True). The return_exceptions=True flag matters: cancelled tasks would otherwise raise CancelledError out of the gather and abort the rest of your cleanup.
import asyncio tasks = set() async def job(name): try: while True: await asyncio.sleep(0.01) except asyncio.CancelledError: print(name, 'cancelled') raise async def main(): for name in ('sync', 'report'): t = asyncio.create_task(job(name)) tasks.add(t) t.add_done_callback(tasks.discard) await asyncio.sleep(0.05) pending = list(tasks) for t in pending: t.cancel() results = await asyncio.gather(*pending, return_exceptions=True) print([type(r).__name__ for r in results]) print('still tracked:', len(tasks)) asyncio.run(main())
sync cancelled report cancelled ['CancelledError', 'CancelledError'] still tracked: 0
If you run your own loop instead of uvicorn or aiohttp, you can catch the signal yourself with loop.add_signal_handler(signal.SIGTERM, stop.set) (Unix only) and await stop.wait() before running the same drain, cancel and close steps. uvicorn already does this for you.
Setting an app drain timeout of 30 s or more while Kubernetes allows 30 s. SIGKILL arrives first and cuts requests off mid-flight.
Never closing the ClientSession or AsyncClient. You get Unclosed client session warnings and leaked sockets.
Untracked fire-and-forget tasks are never cancelled, so shutdown can hang waiting on them.
Blocking I/O goes to to_thread or your own pool, and CPU work goes to processes. Shut down in the order stop, drain, cancel, close, exit, and keep the whole budget below the Kubernetes grace period.
Part 11 · Async HTTP Cheatsheet
Clients, timeouts and pools
Everything in this chapter comes down to a few habits. Open one client per application, give every call a timeout, size the pool deliberately, and always leave the response block so the connection goes back. The table lines up the two libraries side by side.
| Knob | aiohttp | httpx |
|---|---|---|
| One client | ClientSession opened in cleanup_ctx | AsyncClient opened in lifespan |
| Close | await session.close() | await client.aclose() |
| Timeouts | ClientTimeout(total=10, sock_connect=3) | Timeout(10, connect=3) |
| Pool | TCPConnector(limit=100, limit_per_host=10) | Limits(max_connections=100, max_keepalive_connections=20) |
| Stream | r.content.iter_chunked() | c.stream() + aiter_bytes() |
| Errors | ClientError / ClientResponseError | RequestError / HTTPStatusError / TimeoutException |
The client is created at startup and closed at shutdown. In FastAPI that is the lifespan function, and in aiohttp.web it is a cleanup_ctx generator. Code before the yield runs at startup and code after it runs at shutdown. Handlers only borrow the client.
- 1Startupopen client in lifespan / cleanup_ctx
- 2Handle requestsevery request reuses the pool
- 3Shutdownclose client once after yield
Timeouts need two layers. The library timeout bounds each phase, such as connect and read. A hard total cap is separate. aiohttp has a total field, while httpx has only per-phase limits, so wrap the whole operation in asyncio.timeout() (Python 3.11+). On the deadline it cancels everything inside, including each child of a gather.
# fragment: library calls, not runnable on their own t = aiohttp.ClientTimeout(total=10, sock_connect=3) t = httpx.Timeout(10, connect=3) # startup: c = AsyncClient(...) async with asyncio.timeout(15): # hard total cap async with c.stream('GET', u) as r: async for b in r.aiter_bytes(): f.write(b) # block exit -> connection back to pool # shutdown: await c.aclose()
Phase limits plus a hard deadline around a streamed download
The pool knobs cap sockets, not tasks. limit_per_host=10 protects one upstream from you. If the pool is full, aiohttp makes the request wait and that wait counts toward the connect and total timeouts. httpx raises PoolTimeout instead.
Phase timeouts stop a stalled socket. Only a total deadline stops a slow trickle, so set both.
Fan-out and retries
Calling gather() over 10,000 URLs starts 10,000 tasks at once, which means a memory spike, queued pool waits and a pile of 429 responses. Put a Semaphore(n) around the work so only n coroutines are inside at once. All the tasks still exist, but only n of them are active.
import asyncio active = 0 peak = 0 async def fetch(i): global active, peak active += 1 peak = max(peak, active) await asyncio.sleep(0.01) active -= 1 return i async def main(): sem = asyncio.Semaphore(5) async def bounded(i): async with sem: return await fetch(i) rs = await asyncio.gather(*(bounded(i) for i in range(50))) print(len(rs), peak) asyncio.run(main())
50 tasks, never more than 5 in flight
50 5
| Situation | Tool | Why |
|---|---|---|
| At most n in flight | Semaphore(n) + gather / TaskGroup | simple and bounded |
| Huge or streaming input | Queue(maxsize=...) + N workers | flat memory |
| Requests per second | aiolimiter.AsyncLimiter(10, 1) | a semaphore is not a rate |
| Sockets to one host | limit_per_host | the pool, not your tasks |
Retries are only safe for failures that are both transient and idempotent. That means connect errors, timeouts, and the statuses 429, 502, 503 and 504, on GET, PUT or DELETE. A POST is retried only with an Idempotency-Key header. Wait with exponential backoff and full jitter: sleep a random time between zero and the capped ceiling, so clients don't retry in lockstep.
import random def ceiling(n, base=0.5, cap=30): return min(cap, base * 2**n) def backoff(n): return random.uniform(0, ceiling(n)) for n in (0, 1, 3, 6, 8): print(n, ceiling(n)) print(all(0 <= backoff(n) <= ceiling(n) for n in range(10)))
The ceiling doubles until the cap; the sleep is anywhere below it
0 0.5 1 1.0 3 4.0 6 30 8 30 True
If a 429 or 503 carries a Retry-After header, sleep that long (capped) instead of your computed delay. Sleep after releasing the semaphore so a slot isn't held idle. Always re-raise asyncio.CancelledError, otherwise shutdown and asyncio.timeout() stop working.
Streaming, servers and shutdown
Streaming reads a body chunk by chunk instead of loading it whole. In aiohttp, loop over r.content.iter_chunked(n). In httpx, use c.stream() and then aiter_bytes(), because a plain c.get() buffers the entire body. Call raise_for_status() before iterating so an error page isn't saved as data. Always exit the async with, even if you stop early, to release the connection.
The server rule is about what runs on the event loop. An async def handler runs on the loop, so it may only await non-blocking I/O. Blocking libraries belong in a plain def endpoint, which FastAPI runs in a threadpool, or behind asyncio.to_thread. CPU-bound work goes to a ProcessPoolExecutor, because threads share the GIL.
import asyncio import time def blocking(): time.sleep(0.1) return 'done' async def ticker(ticks): for _ in range(3): await asyncio.sleep(0.02) ticks.append(1) async def main(): ticks = [] r, _ = await asyncio.gather(asyncio.to_thread(blocking), ticker(ticks)) print(r, len(ticks)) asyncio.run(main())
The loop keeps ticking while the blocking call sits in a thread
done 3Shutdown happens in a fixed order and has to finish inside the platform's grace period, for example 30 seconds on Kubernetes. After SIGTERM the grace clock starts. Stop accepting new connections, drain in-flight requests, cancel background tasks, and only then close the clients.
- 1SIGTERMgrace clock starts
- 2Stop acceptingclose the listener
- 3Drain in-flightlet requests finish
- 4Cancel bg taskscancel, then gather
- 5Close clientssession / client
Catch errors by layer. For a failed status use ClientResponseError in aiohttp or HTTPStatusError in httpx. For transport problems use ClientError or RequestError, and for timeouts use TimeoutException in httpx. Status errors and transport errors are different failures, so handle them separately, and map upstream failures to 502 or 504 instead of leaking a raw 500.
Top mistakes
These five cause most real outages in async HTTP code. Each one is cheap to avoid once you know it.
A new ClientSession or AsyncClient for every call means a fresh TCP and TLS handshake each time and no pooling. Open one at startup and share it.
requests.get() or time.sleep() blocks the whole event loop, so every other request stalls. Use an async client, asyncio.sleep, or push the call into to_thread.
gather() over 10,000 URLs starts them all together. Bound it with a Semaphore, a queue with workers, or aiolimiter for a request rate.
A hung server hangs you, and the stuck await also holds a pool slot. Set phase timeouts and wrap the whole operation in asyncio.timeout().
A response that is never read or exited leaks its connection until the pool is empty and later calls hang. Always use async with on responses and streams.
Part 12 · Check yourself
Quiz
Each call to this helper works, and the data comes back correct. What does it cost you, and which variant would produce 'Unclosed client session' warnings?
- Every call builds a new ClientSession, so it also builds a new connector and a new pool. Nothing is ever reused, and each request pays a fresh TCP + TLS handshake (1–3 extra RTTs).
- This code closes its session through
async with, so it does not leak sockets or warn. It only wastes time and pools nothing. - The 'Unclosed client session' warning belongs to a variant like
s = aiohttp.ClientSession()with noasync withand noawait s.close(). - Fix: open one session at startup (cleanup_ctx or lifespan), pass it in, and close it once on shutdown.
async def get_json(url): async with aiohttp.ClientSession() as s: async with s.get(url) as r: r.raise_for_status() return await r.json()
This retry helper is called under a shared semaphore. Name four problems.
- The bare
except:also catchesasyncio.CancelledError, so timeouts and shutdown cannot stop the helper. Catch specific transport errors and re-raise cancellation. asyncio.sleepruns insideasync with sem, so a backing-off task holds a concurrency slot while it does nothing. Sleep after the block exits.- The delay
2 ** nhas no jitter, so many failed tasks wake at the same instant and hit the upstream together. Userandom.uniform(0, delay)with a cap. - Any status other than 200 is retried, including 401, 404 and 422, which will fail again. Retry only 429, 502, 503 and 504, and honor Retry-After.
async def get_retry(c, url): async with sem: for n in range(5): try: r = await c.get(url) if r.status_code == 200: return r except: pass await asyncio.sleep(2 ** n)
Three requests take 1s, 2s and 3s. They are launched with gather() and each one runs inside async with asyncio.Semaphore(2). About how long does the whole batch take, and what would it take without the semaphore or with a plain for loop?
- About 4 seconds. Requests 1 and 2 start at t=0, request 1 finishes at t=1 and frees a slot, and request 3 then runs from t=1 to t=4.
- Without the semaphore, gather() overlaps all three, so the total is about the slowest one: 3 seconds.
- A sequential
for u in urls: await fetch(u)has no concurrency and takes the sum: 6 seconds. - The semaphore trades some speed for a bounded number of requests in flight.
A FastAPI service has 50 clients hitting /slow at once. Why does /health stop answering on the first version, and why does the second version keep working?
- An
async defendpoint runs directly on the event loop.requests.getblocks the only thread, so every coroutine, including/health, waits until it returns. - A plain
defendpoint runs in the anyio threadpool (40 threads by default), so the loop stays free. It is safe, but capped at about 40 concurrent calls. - The better fix is
async defwithawait client.get(...)on the lifespan httpx client, which needs no thread at all. - A symptom to watch for is latency spiking on every endpoint at once. PYTHONASYNCIODEBUG=1 logs callbacks slower than 100ms.
@app.get('/slow') async def slow(): return requests.get(URL).json() @app.get('/slow') def slow(): return requests.get(URL).json()
An httpx client has timeout=httpx.Timeout(5.0). The server sends one byte every 4 seconds for an hour. Does the call raise ReadTimeout, and how do you bound it?
- No. httpx timeouts apply per phase, and each read gets a fresh 5s window. A byte every 4s never trips it.
- aiohttp's
totalfield would cap the whole operation, but httpx has no total. - Put a hard deadline on the caller with
async with asyncio.timeout(10):(3.11+), orasyncio.wait_for(coro, 10)on older versions. - When the deadline hits, the pending request is cancelled and the caller gets TimeoutError, whatever the phase limits say.
async with asyncio.timeout(10): r = await c.get(url)
Summary
- Open one ClientSession or AsyncClient per app in lifespan or cleanup_ctx, reuse it everywhere and close it once on shutdown.
- Always set timeouts: phase limits on the client plus
asyncio.timeout()as a hard total, since httpx has no total. - Bound fan-out with a Semaphore or a Queue with workers, and keep n at or below the per-host pool limit.
- Retry only transient, idempotent failures with capped exponential backoff and full jitter, honor Retry-After, sleep outside the semaphore and re-raise CancelledError.
- Stream large bodies with
r.content.iter_chunked()orc.stream(), and always exit the context so the connection returns to the pool. - In an
async defhandler only await non-blocking I/O; send blocking calls to adefendpoint orasyncio.to_thread, and CPU work to a ProcessPoolExecutor. - Shut down in order: stop accepting, drain, cancel background tasks, close clients, and keep the drain time below the Kubernetes grace period.