Smoothing Bursts with a Leaky Bucket¶
Token buckets and leaky buckets are often described as the same algorithm. Configured to the same average rate, they produce very different traffic. A token bucket saves up unused capacity and spends it in a burst; a leaky bucket releases calls at a fixed interval no matter how they arrive. Measured with 100 concurrent callers at 50 per second: the token bucket (burst 20) finished in 1.60 s but let 24 calls through in a single 100 ms window, with some calls back to back; the leaky bucket took 1.98 s, never released more than 6 in any 100 ms, and kept at least 19 ms between consecutive calls. When the downstream punishes bursts — an API that enforces per-second limits strictly, a device on a serial bus, a database that degrades under spikes — the leaky bucket is the one you want. This guide implements it with bounded capacity and correct cancellation.
Prerequisites¶
- Python 3.11+, stdlib only.
- Token buckets, from token bucket rate limiter for asyncio clients.
- Timer precision, from using loop.call_later and timer handles.
1. Schedule each caller into the next free slot¶
The leaky bucket as a meter is easiest to implement as a slot scheduler: keep the time of the next free slot, give each caller that slot, and advance it by one interval:
import asyncio
import time
class LeakyBucket:
def __init__(self, rate: float, capacity: int) -> None:
self.interval = 1.0 / rate
self.capacity = capacity # how many callers may wait at once
self.next_slot = time.monotonic()
self.waiting = 0
async def acquire(self) -> None:
if self.waiting >= self.capacity:
raise OverflowError("leaky bucket full")
self.waiting += 1
try:
now = time.monotonic()
slot = max(now, self.next_slot)
self.next_slot = slot + self.interval
await asyncio.sleep(slot - now)
finally:
self.waiting -= 1
The slot assignment happens synchronously, before any await, so concurrent callers on one loop get distinct, ordered slots without a lock. An idle bucket does not accumulate credit: max(now, self.next_slot) means a caller arriving after a quiet period is released immediately, and the next one an interval later — never a burst.
Verify: release times of 100 concurrent callers are spaced by interval within timer precision.
2. Compare the traffic it produces¶
Measure the peak rate in a short window, not just the average:
def max_in_window(timestamps: list[float], window: float = 0.1) -> int:
ts = sorted(timestamps)
best = j = 0
for i in range(len(ts)):
while ts[i] - ts[j] > window:
j += 1
best = max(best, i - j + 1)
return best
Results for 100 concurrent callers at 50 per second:
| Limiter | Total time | Max in any 100 ms | Min gap |
|---|---|---|---|
| token bucket, burst 20 | 1.60 s | 24 | 0.00 ms |
| leaky bucket | 1.98 s | 6 | 19.09 ms |
The token bucket finished sooner because its first 20 calls cost nothing; the leaky bucket paid one interval per call from the start. Twenty-four calls in 100 ms is a momentary rate of 240 per second — five times the configured rate. If the downstream counts requests per second strictly, those bursts are exactly what earns 429s, as discussed in handling 429 Retry-After responses in async clients.
Verify: compute the peak window rate for your limiter against the downstream's documented limit window.
3. Bound how many callers may wait¶
Without a bound, a leaky bucket under sustained overload queues callers indefinitely: each new caller's slot is further in the future, and latency grows without limit. capacity caps the queue; callers beyond it are rejected immediately:
lb = LeakyBucket(rate=50, capacity=10)
results = await asyncio.gather(*(lb.acquire() for _ in range(30)), return_exceptions=True)
rejected = sum(isinstance(r, OverflowError) for r in results) # 20
Measured: with capacity 10 and 30 simultaneous arrivals, 20 were rejected at once. A capacity of rate × max_acceptable_delay is a natural choice — 50 per second with a 200 ms tolerance gives 10 — because anything beyond it would wait longer than the caller is willing to. Rejection is the leaky bucket's form of backpressure; propagate it to the caller rather than retrying in a loop, which only re-enters the queue.
Verify: under overload, the wait time of admitted callers never exceeds capacity / rate.
4. Handle cancellation without leaking slots¶
A caller cancelled while waiting has already consumed a slot: next_slot moved forward when it was admitted. With many cancellations — timeouts upstream, disconnected clients — the bucket releases fewer calls than its rate, because cancelled slots pass unused. For strict output rates that is harmless; for throughput it is a loss. Reclaiming the slot requires knowing it was the latest one:
async def acquire(self) -> None:
if self.waiting >= self.capacity:
raise OverflowError("leaky bucket full")
self.waiting += 1
now = time.monotonic()
slot = max(now, self.next_slot)
self.next_slot = slot + self.interval
try:
await asyncio.sleep(slot - now)
except asyncio.CancelledError:
if self.next_slot == slot + self.interval: # nobody took a later slot
self.next_slot = slot # give this one back
raise
finally:
self.waiting -= 1
Only the most recently assigned slot can be given back without reordering everyone behind it; earlier slots stay consumed. That keeps the guarantee — never two calls closer than one interval — while recovering the common case of a single timed-out caller. Cancellation semantics for waiting primitives are covered in Cancellation Patterns.
Verify: cancel the last waiter; the next caller gets its slot, and spacing between released calls never drops below the interval.
5. Use it per destination, not per process¶
A leaky bucket shapes traffic toward one destination. Give each downstream — each API host, each tenant's quota, each device — its own bucket, and share it across all tasks in the process that call that destination:
from collections import defaultdict
buckets: dict[str, LeakyBucket] = defaultdict(lambda: LeakyBucket(rate=10, capacity=20))
async def call_partner(host: str, payload: dict):
await buckets[host].acquire()
async with session.post(f"https://{host}/ingest", json=payload) as r:
return r.status
As with any in-process limiter, several worker processes multiply the rate: four workers each at 10 per second send 40 per second to the partner. Divide the rate by the number of processes, or coordinate through Redis — a shared next_slot in a Lua script works the same way — when the downstream's limit is strict, as in enforcing multiple rate limits at once.
Verify: the partner's per-second request count, measured on their side or from your access logs, never exceeds the agreed rate.
Verification¶
The leaky bucket is correct when:
- Released calls are spaced by at least one interval, measured.
- Peak window rate stays at the configured rate, unlike a token bucket's bursts.
- Waiting callers are bounded by
capacity, and overflow is rejected immediately. - Cancelled waiters do not leave permanent gaps where reclaimable.
Diagnostic Hook: export the number of waiting callers, rejections per minute, and the current scheduling delay (next_slot - now). A delay that sits near capacity / rate means the bucket is saturated — the downstream rate is lower than demand — and rejections will follow; that is the signal to negotiate a higher limit or shed lower-priority calls before they queue.
Pitfalls & edge cases¶
- Unbounded capacity. Latency grows without limit under overload.
- Accumulating credit while idle. That is a token bucket; a leaky bucket must not burst after a quiet period.
- Per-process buckets for a global limit. Multiply by the number of processes.
- Expecting sub-millisecond spacing.
asyncio.sleepprecision is around a millisecond; very high rates need batching instead.
Frequently Asked Questions¶
What is the difference between a token bucket and a leaky bucket?
A token bucket stores unused capacity and spends it in bursts; a leaky bucket releases calls at a fixed interval regardless of arrivals. At 50 calls per second, a token bucket with burst 20 allowed 24 calls in 100 ms, a leaky bucket 6.
How do I implement a leaky bucket in asyncio?
Keep the time of the next free slot. Each caller takes max(now, next_slot), advances next_slot by one interval, and sleeps until its slot. Do the assignment before any await so concurrent callers get distinct slots.
When should I use a leaky bucket instead of a token bucket?
When the destination enforces strict per-second limits or degrades under spikes, so that the momentary rate must never exceed the average. Use a token bucket when short bursts are fine and lower latency matters.
What happens when a leaky bucket is overloaded?
Without a bound, waiting time grows indefinitely. Cap the number of waiting callers, for example at rate multiplied by the acceptable delay, and reject the rest immediately.
Related¶
- Rate Limiting & Throttling — up to the topic overview.
- Adaptive concurrency limits with AIMD — when the safe rate is unknown and must be discovered.
- Concurrent Execution & Worker Patterns — the section overview.