Resilience — Semaphores & Rate Limits¶
asyncio.gather() over 10,000 URLs opens 10,000 requests at once: the target answers 429, the connection pool times out, and your own process runs out of sockets. Two different limits prevent that, and they are easy to mix up:
| Concurrency limit | Rate limit | |
|---|---|---|
| Limits | How many calls are in flight at the same time | How many calls start per time window |
| Typical rule | "At most 10 open requests to this host" | "At most 100 requests per minute per API key" |
| Python tools | asyncio.Semaphore, threading.Semaphore, anyio.CapacityLimiter, ThreadPoolExecutor(max_workers), httpx.Limits |
aiolimiter.AsyncLimiter; Redis counters across processes |
| Protects | Your process (sockets, memory, threads) and the target's workers | The target's quota; avoids 429 |
| Fast calls | 10 slots × 20 ms calls = up to 500 calls/s | 100/min no matter how fast the calls are |
Most clients need both: a semaphore for in-flight calls and a rate limiter for the provider's quota.
asyncio: Semaphore and TaskGroup¶
# app/limits.py
import asyncio
import httpx
from tenacity import retry, retry_if_exception, stop_after_attempt, wait_random_exponential
from app.http_client import is_retryable
async def fetch_all(client: httpx.AsyncClient, urls: list[str], *, max_in_flight: int = 10) -> list[dict]:
"""GET every URL, at most `max_in_flight` requests at the same time."""
semaphore = asyncio.Semaphore(max_in_flight)
async def fetch(url: str) -> dict:
async with semaphore: # waits here while the limit is reached
response = await client.get(url)
response.raise_for_status()
return response.json()
async with asyncio.TaskGroup() as tg: # one failure cancels the rest
tasks = [tg.create_task(fetch(url)) for url in urls]
return [task.result() for task in tasks]
@retry(
retry=retry_if_exception(is_retryable),
stop=stop_after_attempt(3),
wait=wait_random_exponential(multiplier=0.2, max=5),
reraise=True,
)
async def get_json(client: httpx.AsyncClient, semaphore: asyncio.Semaphore, url: str) -> dict:
async with semaphore: # the slot is free while tenacity sleeps
response = await client.get(url)
response.raise_for_status()
return response.json()
fetch_all against a mocked API (50 URLs, max_in_flight=5, each response takes 50 ms) finished in 0.53 s with at most 5 requests in flight — ten waves of five. get_json is used in Limits and Retries Together below.
- Create the semaphore inside the running loop for the work that shares it. One semaphore per batch (as here) limits the batch; one per client or per host limits the whole process.
asyncio.TaskGroup(Python 3.11+) waits for all tasks and cancels the rest when one fails;asyncio.gather(..., return_exceptions=True)collects all results and errors instead.- All tasks are created up front. That is fine for thousands of items. For millions, start N worker tasks that read from an
asyncio.Queue(05) — the queue size becomes the limit and memory stays flat. BoundedSemaphoreraisesValueError: BoundedSemaphore released too many timeson an extrarelease(). A plainSemaphoresilently grows (Semaphore(1)released twice allows 2 at once). Useasync with, and preferBoundedSemaphorewhen you callacquire()/release()by hand.
Threads: Semaphore and ThreadPoolExecutor¶
import threading
import time
from concurrent.futures import ThreadPoolExecutor
license_slots = threading.BoundedSemaphore(2) # e.g. a device farm with 2 free phones
def run_on_device(test: str) -> str:
with license_slots: # blocks while 2 tests hold a device
time.sleep(0.1)
return f"{test}: passed"
start = time.perf_counter()
with ThreadPoolExecutor(max_workers=8) as pool: # 8 threads, but only 2 inside at once
results = list(pool.map(run_on_device, [f"test_{i}" for i in range(6)]))
print(results[:2], f"{time.perf_counter() - start:.1f} s")
# ['test_0: passed', 'test_1: passed'] 0.3 s
ThreadPoolExecutor(max_workers=n)is a concurrency limit: at mostnjobs run, the rest wait in its internal queue. The default ismin(32, os.process_cpu_count() + 4)on Python 3.13+ (os.cpu_count()before).- A semaphore inside the pool limits one resource (devices, licences, a fragile API) while the pool runs other work in parallel.
threading.BoundedSemaphoreraisesValueErroron an extra release, like the asyncio one.
anyio.CapacityLimiter¶
anyio runs on asyncio and Trio. Its CapacityLimiter is a semaphore with extras — and the limit for blocking code sent to threads:
import time
import anyio
import anyio.to_thread
db_slots = anyio.CapacityLimiter(3) # at most 3 blocking DB calls at once
def blocking_query(n: int) -> int:
time.sleep(0.1) # a sync driver, a file read, a CPU-light C call
return n * n
async def query(n: int, results: list[int]) -> None:
results.append(await anyio.to_thread.run_sync(blocking_query, n, limiter=db_slots))
async def main() -> None:
results: list[int] = []
start = time.perf_counter()
async with anyio.create_task_group() as tg:
for n in range(9):
tg.start_soon(query, n, results)
print(sorted(results), f"{time.perf_counter() - start:.1f} s") # 9 calls / 3 slots * 0.1 s
print(db_slots.total_tokens, db_slots.borrowed_tokens)
print(anyio.to_thread.current_default_thread_limiter().total_tokens) # default for run_sync: 40
anyio.run(main)
# [0, 1, 4, 9, 16, 25, 36, 49, 64] 0.3 s
# 3 0
# 40
total_tokenscan be changed at runtime (scale a limit up or down without recreating it);borrowed_tokensandstatistics()show current use.anyio.to_thread.run_sync()withoutlimiter=shares one default limiter of 40 threads — FastAPI / Starlette run sync endpoints and dependencies through it, so 40 slow sync endpoints block the 41st request.
Rate Limits with aiolimiter¶
import asyncio
import time
from aiolimiter import AsyncLimiter
async def main() -> None:
limiter = AsyncLimiter(5, 1) # 5 per 1 second (default period: 60 s!)
start = time.perf_counter()
stamps: list[float] = []
async def call() -> None:
async with limiter:
stamps.append(round(time.perf_counter() - start, 2))
async with asyncio.TaskGroup() as tg:
for _ in range(20):
tg.create_task(call())
print(stamps[:7], "...", stamps[-1])
asyncio.run(main())
# [0.0, 0.0, 0.0, 0.0, 0.0, 0.2, 0.4] ... 3.0
AsyncLimiter(max_rate, time_period=60): the period defaults to 60 seconds —AsyncLimiter(10)means 10 per minute, not per second.- It is a leaky bucket: up to
max_rateacquisitions pass at once (a burst), then one everytime_period / max_rateseconds. If the provider forbids bursts, use a smallermax_ratewith a proportionally smaller period (AsyncLimiter(1, 0.2)instead ofAsyncLimiter(5, 1)). has_capacity()checks without waiting;acquire(amount)takes several units (e.g. tokens for an LLM call).- One limiter belongs to one event loop. Reusing it across
asyncio.run()calls (common in tests) emitsRuntimeWarning: This AsyncLimiter instance is being re-used across loops; create it per loop or per test. - For sync code or limits shared across processes see
pyrate-limiterandlimits, or a Redis counter (Redis — Patterns).
Limits per Host¶
A crawler, a test-data loader or a contract-test runner talks to many hosts. One global semaphore lets a slow host take all slots; one per host keeps them independent:
# app/host_limits.py
import asyncio
from collections import defaultdict
from urllib.parse import urlsplit
import httpx
from aiolimiter import AsyncLimiter
class HostLimits:
"""Concurrency limit per host plus a requests-per-second limit per host."""
def __init__(self, max_in_flight: int = 4, max_rate: float = 10, period: float = 1.0) -> None:
self._semaphores: defaultdict[str, asyncio.Semaphore] = defaultdict(
lambda: asyncio.Semaphore(max_in_flight)
)
self._rates: defaultdict[str, AsyncLimiter] = defaultdict(
lambda: AsyncLimiter(max_rate, period)
)
async def get(self, client: httpx.AsyncClient, url: str) -> httpx.Response:
host = urlsplit(url).netloc
async with self._semaphores[host]: # outer: limit requests in flight
async with self._rates[host]: # inner: send right after getting a token
return await client.get(url)
With max_in_flight=2, ten requests to a.test and ten to b.test ran with a peak of 4 in flight — 2 per host.
- Semaphore outside, rate limiter inside: a task takes a token only when it can send right away. The other order lets tasks collect tokens while they wait for a slot and then send them in a burst.
- HTTPX has its own pool limits per client —
httpx.Limits(max_connections=100, max_keepalive_connections=20)by default — for all hosts together. Above the limit, requests wait for a connection and fail withhttpx.PoolTimeoutafter thepooltimeout (5 s by default).
Limits and Retries Together¶
A retry loop inside a semaphore keeps the slot while it sleeps between attempts. Twenty jobs, a limit of 4, each job's first call fails and waits 0.5 s before the retry:
import asyncio
import time
from tenacity import retry, retry_if_exception_type, stop_after_attempt, wait_fixed
class Busy(Exception):
pass
failed_once: set[int] = set()
async def call_api(job: int) -> int:
await asyncio.sleep(0.1) # the request itself
if job not in failed_once: # every job fails once, then succeeds
failed_once.add(job)
raise Busy(job)
return job
retry_busy = retry(retry=retry_if_exception_type(Busy), stop=stop_after_attempt(3), wait=wait_fixed(0.5))
@retry_busy
async def call_with_retry(job: int) -> int:
return await call_api(job)
async def hold_slot(sem: asyncio.Semaphore, job: int) -> int:
async with sem: # slot is held during the 0.5 s backoff sleep
return await call_with_retry(job)
@retry_busy
async def release_slot(sem: asyncio.Semaphore, job: int) -> int:
async with sem: # slot is held only while the request runs
return await call_api(job)
async def run(strategy) -> float:
failed_once.clear()
sem = asyncio.Semaphore(4)
start = time.perf_counter()
async with asyncio.TaskGroup() as tg:
for job in range(20):
tg.create_task(strategy(sem, job))
return time.perf_counter() - start
print(f"hold slot during backoff: {asyncio.run(run(hold_slot)):.2f} s")
print(f"release slot during backoff: {asyncio.run(run(release_slot)):.2f} s")
hold slot during backoff: 3.52 s
release slot during backoff: 1.11 s
Holding the slot wastes it on sleeping: three times slower here, and under a real outage all slots end up sleeping while healthy work waits. Put the semaphore inside the retried function (like get_json above), so every attempt takes a slot and gives it back before the backoff sleep.
The trade-off: a released slot goes to the next waiting task, so a retrying task queues again behind new work. If retries must finish first (e.g. they hold a lock elsewhere), keep the slot — and keep the backoff short.
Rate limiters and retries:
- Every retry consumes a token — retries count against the provider's quota exactly like first attempts.
- A
429means the limiter is too generous (or other clients share the quota). HonourRetry-After, then lower the rate rather than retrying harder. - Shared quota, shared limiter: all tasks that use one API key go through one limiter instance.
Distributed Limits¶
Semaphores and AsyncLimiter live in one process. Four pods with Semaphore(10) allow 40 concurrent calls; eight pytest-xdist workers with AsyncLimiter(5, 1) send 40 requests per second. Options:
- Divide the limit by the number of instances (simple; wrong when instances scale up or down).
- A shared counter in Redis — fixed or sliding window, or a token bucket in a Lua script: Redis — Patterns: Rate Limiting.
- Let the provider enforce it: honour
429andRetry-After, back off, and keep a local limiter slightly below the quota. - One gateway in front of the dependency (API gateway, LiteLLM proxy for LLM calls) that applies the limit for every client.
Checklist¶
- Every fan-out (
gather,TaskGroup, thread pool) has an explicit concurrency limit - Per-host or per-dependency limits where one slow host must not block others
- Rate limits use
AsyncLimiter(rate, period)with an explicit period - Semaphore outside, rate limiter inside; the semaphore is released during retry backoff
- Retries and
429responses are counted against the same limiter - Limits that must hold across processes or pods live in Redis or a gateway, not in a
Semaphore - Tests measure peak concurrency and assert it equals the limit (06)
See also¶
- Resilience — Timeouts, Fallbacks & Circuit Breakers
- Resilience — Race Conditions
- HTTPX — Async Patterns
- Redis — Patterns: Caching, Rate Limits, Locks, Messaging
- Context Managers & Async