Rate-Limited Worker Pools in asyncio¶
A worker pool that calls a rate-limited API needs two numbers to agree: the rate limit, which caps how often calls start, and the worker count, which caps how many run at once. Get the relationship wrong and the pool either never reaches the limit or exceeds it in bursts. This guide uses a simulation — workers call asyncio.sleep in place of an API, so the numbers isolate the pool's behaviour — on Python 3.14, with a token bucket set to 100 calls per second and 1,000 items. With calls taking a fixed 100 ms, 5 workers managed only 49.8 starts per second; 10 workers reached 99.1, and more added nothing. With latencies spread like a real API (median 71 ms, p99 461 ms), 10 workers reached 89.3 per second and 20 were needed for 98.1. A bucket burst of 100 let 109 calls start within 100 ms, against 10 with a burst of 1. And when the upstream also capped concurrent calls at 8, taking the rate token before the concurrency slot produced spikes of 18 starts in 100 ms; taking the slot first never exceeded 10, at the cost of throughput — 77.4 against 85.1 per second.
Prerequisites¶
- A token bucket, from a token bucket rate limiter for asyncio clients.
- A worker pool, from building an async worker pool with TaskGroup.
- The topic overview, Worker Pool Implementations.
1. Share one limiter across all workers¶
The rate limit belongs to the upstream, so one limiter instance must be shared by every worker that calls it — not one per worker, which would multiply the rate by the worker count:
async def run_pool(items, call, rate=100, workers=20):
limiter = TokenBucket(rate=rate, burst=1) # one instance for the pool
q = asyncio.Queue()
for item in items:
q.put_nowait(item)
async def worker():
while not q.empty():
item = q.get_nowait()
await limiter.acquire() # wait for permission to start
await call(item)
async with asyncio.TaskGroup() as tg:
for _ in range(workers):
tg.create_task(worker())
Acquire the permit immediately before the call, after taking the item, so time spent waiting for the limiter is not counted against any per-call timeout. Record the time each call starts; every measurement below comes from those timestamps.
Verify: count call starts per second from recorded timestamps; with enough workers, the count matches the configured rate.
2. Size workers by rate times latency¶
A worker that is busy with a call cannot start another, so the pool can start at most workers / latency calls per second. To reach a rate, it needs at least rate × latency workers — Little's law applied to the pool:
def workers_needed(rate_per_s: float, latency_s: float, headroom: float = 2.0) -> int:
return math.ceil(rate_per_s * latency_s * headroom)
Measured with a fixed 100 ms call and a 100/s limit: 5 workers started 49.8 calls per second and took 20.2 s for 1,000 items; 10 workers reached 99.1 per second and 10.2 s; 12, 20 and 50 workers were no faster. The limiter was in control once there were enough workers to keep it busy. The headroom factor matters because real latency varies, as the next step measures.
Verify: compute rate × mean latency from production metrics, and check that the worker count is above it with margin.
3. Account for latency variance¶
Real calls do not all take the mean. With latencies drawn from a log-normal distribution — median 71 ms, mean 97 ms, p99 461 ms — the pool needed more workers than the mean suggests. 10 workers, enough for fixed 100 ms calls, reached only 89.3 per second, because at any moment a few workers were stuck in slow calls. 12 workers reached 95.7, and 20 reached 98.1. A factor of two over rate × mean latency was enough here; a heavier tail needs more.
Extra workers cost little when the limiter is the bottleneck: they wait inside acquire(), holding no connections. They do raise how many calls can be in flight at once, which matters if the upstream also limits concurrency — the subject of step 5.
Verify: the pool reaches the target rate in a test with production's latency distribution, not just its mean.
4. Keep the bucket's burst small¶
A token bucket's burst size lets that many calls start at once after an idle period. A worker pool with a backlog is never idle at the start, so the whole burst is spent immediately:
limiter = TokenBucket(rate=100, burst=1) # smooth: one start every 10 ms
Measured with 200 workers and a full queue: burst 1 allowed at most 10 starts in any 100 ms window, which is exactly the rate. Burst 20 allowed 29, and burst 100 allowed 109 — the whole burst plus the window's refill, in the first tenth of a second. Upstream limiters that count requests in short windows reject that spike with 429 responses even though the long-run average is within the limit; the measured average was 111 per second over the run with burst 100, against 98.6 with burst 1. Use a burst of 1 for batch pools, and only raise it when the upstream documents a burst allowance.
Verify: the maximum number of call starts in any 100 ms window stays at or below rate × 0.1 plus the configured burst.
5. Combine a rate limit with a concurrency limit¶
Many APIs limit both how often calls start and how many may run at once. With a semaphore for the second, the order of the two waits matters:
async def worker(): # slot first, then token
while not q.empty():
item = q.get_nowait()
async with concurrency: # asyncio.Semaphore(8)
await limiter.acquire()
await call(item)
Measured with a 100/s rate, a concurrency limit of 8, 50 workers and the heavy-tailed latencies: taking the rate token first and then waiting for a slot let tokens accumulate in waiting workers, and when slots freed up, several calls started together — up to 18 in a 100 ms window, almost twice the rate. Taking the slot first and then the token never exceeded 10 starts in 100 ms. The cost: a worker holds a slot while it waits for a token, so throughput fell from 85.1 to 77.4 starts per second. When the upstream counts requests in short windows, slot-then-token avoids rejections; when it only enforces a long-run average, token-then-slot is faster.
Verify: with both limits configured, record starts and in-flight counts together, and check each against its own limit.
Verification¶
A rate-limited pool is configured correctly when:
- One limiter is shared by every worker calling the same upstream.
- Worker count exceeds
rate × latencywith margin for the latency tail — 20 workers for 100/s at a 97 ms mean here. - The burst is 1 unless the upstream allows more, and the busiest 100 ms window stays within the limit.
- Combined limits are acquired slot first when the upstream counts short windows.
Diagnostic Hook: when a pool runs below its configured rate, compare the worker count with rate × observed mean latency. At 5 workers and 100 ms calls, the pool started 49.8 calls per second — exactly workers / latency — so the limiter was never the constraint.
Pitfalls & edge cases¶
- One limiter per worker. The effective rate becomes rate × workers.
- Too few workers. Measured: 5 workers reached 49.8/s of a 100/s limit.
- A large burst on a backlog. Measured: 109 starts in the first 100 ms.
- Token before slot under a concurrency cap. Measured: 18 starts in 100 ms against 10 allowed.
Frequently Asked Questions¶
How many workers does a rate-limited pool need?
At least rate × mean latency, plus margin for slow calls. At 100/s with 100 ms calls, 10 sufficed for fixed latency; with a long tail, 20 were needed to reach 98/s.
Should each worker have its own rate limiter?
No. The limit belongs to the upstream, so all workers must share one limiter; per-worker limiters multiply the effective rate by the number of workers.
Why does my pool get 429 errors at the start of a run?
The token bucket's burst is spent immediately on the backlog. A burst of 100 let 109 calls start in 100 ms; a burst of 1 kept it to 10.
Rate limit or concurrency limit first?
Take the concurrency slot first, then the rate token, if the upstream counts short windows: it stayed within 10 per 100 ms, against 18 the other way, at 77.4/s versus 85.1/s.
Related¶
- Worker Pool Implementations — up to the topic overview.
- Collecting errors from worker pools — handling the calls that still fail.
- Concurrent Execution & Worker Patterns — the section overview.