System design
Part 4 of 6 · Rate limitingDistributed Rate Limits Across Gateways
A limit of 100 req/s per key means nothing if each of N gateway pods enforces it on its own: the real limit becomes N x 100 and changes every time the fleet autoscales. Accurate limits need one shared counter (Redis with an atomic Lua script); fast limits need a local check that costs no network hop. Production systems layer them: the edge drops obvious abuse, a local per-pod bucket rejects clear overage for free, and a global Redis bucket enforces the real quota. When Redis is slow or down, degrade to the local share and alert, instead of either blocking everything or admitting everything.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Overview
A limit of 100 req/s per key means nothing if each of N gateway pods enforces it on its own: the real limit becomes N x 100 and changes every time the fleet autoscales. Accurate limits need one shared counter. Fast limits need a local check that costs no network hop.
By the end you should be able to:
- Show why a correct local limiter still drifts at fleet scale
- Draw edge, local bucket, then global Redis
- Size a hybrid front so most rejects never touch Redis
- Pick regional limits plus an async global quota
- Say what happens when the shared store times out
Why it matters
Interview signal. "Add a rate limiter" is easy on one box and wrong on twenty. You are expected to explain why local limits multiply with pod count, and to pick local, global, or hybrid with an accuracy versus latency argument, including what happens when the shared store fails.
Production signal. Paid quotas silently grow after a scale-out. A Redis hot key adds p99 to every request. Or the limiter itself takes the API down because it failed closed.
Core concepts (deep)
Why local-only limits drift. A load balancer spreads one client across all pods, so each pod sees about 1/N of that client. If every pod allows the full limit, the client gets up to N x limit. If every pod allows limit/N, the total is right only when the balancer is perfectly even. Sticky connections, HTTP/2 multiplexing, and uneven pod counts per zone break that, and the limit still changes on every scale event.
Global limits. One counter per key in a shared store. Redis is the usual choice: a short Lua script (EVALSHA) refills and debits in one atomic step, so two pods cannot both spend the last token. See Redis + Lua. Cost: one network round trip per request (often 0.2-1 ms in the same zone) and a hot key when one tenant is very busy. On Redis Cluster, give every key for one identity the same hash tag so a script never touches two slots.
Hybrid: local front plus global truth. Each pod keeps a small local bucket sized to its share with headroom (for example 1.5 x limit/N). Requests the local bucket rejects never reach Redis, which removes most store traffic during an abusive burst. Requests that pass still debit the global bucket, so the global number stays exact. A variant is batching: each pod leases a block of tokens from Redis (say 10 at a time) and spends them locally, trading a small overshoot for far fewer round trips.
Edge, gateway, service. Put the cheapest checks first. A CDN or WAF rule per IP and path stops scrapers and floods before TLS and auth cost anything. The gateway knows the authenticated identity, so per-user and per-tenant quotas live there. Services can add their own concurrency limits to protect a database. Each layer answers a different question, so they complement rather than replace each other.
Multi-region. One global Redis across regions puts a cross-region round trip (tens of milliseconds) on every request and makes the limiter depend on the WAN. The usual design keeps strong limits per region and reconciles a global quota asynchronously: regions report usage every few seconds and receive a fresh allocation. The cost is a bounded overshoot during the sync interval, which is fine for abuse protection and for most billing quotas.
Failure modes of the store. Fail-open (admit when Redis times out) keeps the API up but removes protection. Fail-closed (reject) turns a Redis blip into a full outage. The common middle path is fail-open to the local per-pod share with a short timeout (a few ms), a circuit breaker so pods stop waiting on a dead store, and an alert. Money-moving or abuse-sensitive endpoints may choose fail-closed explicitly.
Wrong pick, both ways. Local-only for paid or contractual quotas: customers get more than they paid for and the number moves with autoscaling. Global-only for every request across regions: you pay WAN latency on the hot path and the limiter becomes a single point of failure.
40 pods and a sold quota of 100 req/s
Prefer
Local front plus one Redis Lua counter
The local bucket rejects clear overage with no network hop. The sold number lives in one atomic store. A Redis timeout falls back to limit/N and pages someone.
- Local share is about 1.5 x limit/N, not a second copy of the full limit.
- Global EVALSHA is the invoice. Two pods cannot both spend the last token.
- Edge still drops obvious IP and path abuse before auth.
Alternative
Full limit on every pod, or one global Redis across regions
Local-only quotas grow with autoscaling. A single cross-region Redis puts WAN latency on every request and takes the API down with the store.
- 20 pods at 100 req/s is 2000 req/s, and HPA makes it worse.
- limit/N is exact only when the balancer is perfectly even.
- No timeout on the Redis call adds store latency to every request.
Happy path: edge, local bucket, then global Redis
The six steps match Diagram 1. Steps 4a/4b and 6a/6b are the two branches inside steps 4 and 6.
- 1
Edge drops obvious abuse
The CDN or WAF applies cheap per-IP and per-path rules before the request reaches a gateway pod. - 2
Gateway checks the local token bucket
The pod authenticates the caller, builds the key (for example tenant:42), and checks its local bucket. - 3
Local tokens decide the next hop
No local tokens goes to 4a. Tokens left goes to 4b. - 4
Fast 429 or atomic Lua debit
4a replies 429 from the pod with no Redis call. 4b runs the atomic Lua debit on the shared Redis bucket. - 5
Redis answers the global quota
The script says whether the key still has tokens. Nothing else interleaves with that debit. - 6
Forward, 429, or local fallback
6a forwards to the service. 6b replies 429 with Retry-After. If Redis times out, fall back to the strict local share (limit/N) and alert.
Diagrams - step by step
Three small diagrams for distributed limits across gateways. Step numbers in the labels give the animation order. The lesson map under Diagram 1 plays those steps.
Diagram 1 - Happy path: edge, local bucket, then global Redis
Decisions
- 1
Step 1 Edge drops obvious abuse per IP
- nextStep 2 Gateway pod checks a local token bucket
- 2
Step 2 Gateway pod checks a local token bucket
- nextStep 3 Local bucket has tokens?
- ?
Step 3 Local bucket has tokens?
- noStep 4a Fast 429, no Redis call
- yesStep 4b Atomic Lua debit on shared Redis
- 4
Step 4a Fast 429, no Redis call
- 5
Step 4b Atomic Lua debit on shared Redis
- nextStep 5 Global quota left?
- Redis timeoutFailure path - fall back to the local limit only and alert
- ?
Step 5 Global quota left?
- yesStep 6a Forward to the service
- noStep 6b 429 with Retry-After
- 7
Step 6a Forward to the service
- 8
Step 6b 429 with Retry-After
- 9
Failure path - fall back to the local limit only and alert
Local limiters are fast but blind to other pods. The shared Redis store is accurate but adds a network hop. Using both gives low latency for obvious rejects and accuracy for the real quota.
Lesson map
Distributed Rate Limits Across Gateways
Diagram 1 walks 8 steps from Step 1 Edge drops obvious abuse per IP through Step 6b 429 with Retry-After.
Architecture. Step 1 Edge drops obvious abuse per IP Ready. Step 2 Gateway pod checks a local token bucket Ready. Step 3 Local bucket has tokens? Ready. Step 4a Fast 429, no Redis call Ready. Step 4b Atomic Lua debit on shared Redis Ready. Step 5 Global quota left? Ready. Step 6a Forward to the service Ready. Step 6b 429 with Retry-After Ready. Failure path - fall back to the local limit only and alert Ready
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB A["Step 1 Edge drops obvious abuse per IP Ready"] B["Step 2 Gateway pod checks a local token bucket Ready"] C["Step 3 Local bucket has tokens? Ready"] R["Step 4a Fast 429, no Redis call Ready"] D["Step 4b Atomic Lua debit on shared Redis Ready"] E["Step 5 Global quota left? Ready"] F["Step 6a Forward to the service Ready"] R2["Step 6b 429 with Retry-After Ready"] X["Failure path - fall back to the local limit only and alert Ready"] A -->|continues| B B -->|continues| C C -->|no| R C -->|yes| D D -->|continues| E E -->|yes| F E -->|no| R2 D -->|Redis timeout| X
Diagram 2 - Failure path: local-only limits multiply with pod count
Sequence
- 1
Gateway pod 1
Step 1 limit is 100 per s per key, each pod enforces it locally
- 2
Client → Gateway pod 1
Step 2 100 requests, all allowed
- 3
Client → Gateway pod 2
Step 3 100 requests, all allowed
- 4
Client → Gateway pod 3
Step 4 100 requests, all allowed
- 5
Client
Step 5 300 per s reach the service - N pods means N times the limit
- 6
Client
Fix - shared Redis counter, or give each pod limit divided by N
Each pod only sees its own share of traffic, so the effective limit grows with the number of pods and changes every time the fleet scales.
Diagram 3 - Decision: local, global or both
Decisions
- ?
Step 1 Is a rough per-pod cap enough?
- yesLocal per-pod limiter
- noStep 2 Need low latency and an exact key limit?
- 2
Local per-pod limiter
- Wrong pick for paid quotasEffective limit scales with pod count
- ?
Step 2 Need low latency and an exact key limit?
- yesLocal front plus global Redis
- noGlobal Redis with Lua
- 4
Local front plus global Redis
- nextStep 3 Multi-region?
- 5
Global Redis with Lua
- nextStep 3 Multi-region?
- Wrong pick across regionsCross-region latency on every request
- ?
Step 3 Multi-region?
- yesStrong regional limits plus async global quota
- noOne Redis cluster with hash-tagged keys
- 7
Strong regional limits plus async global quota
- 8
One Redis cluster with hash-tagged keys
- 9
Effective limit scales with pod count
- 10
Cross-region latency on every request
Use local limits for coarse protection, global Redis when the number matters, and both when you need speed too. Across regions, keep strong limits regional and reconcile global quotas asynchronously. A local-only limiter is the wrong pick for a paid quota. One global Redis across regions puts cross-region latency on every request.
Working Python
Run python3 distributed_limits.py (stdlib only, under 1 second). SharedStore stands in for Redis plus the Lua token bucket: one lock makes check-and-debit atomic, the same as one EVALSHA call. The in-page sandbox has no operating-system threads, so run this file locally.
"""Distributed rate limits across gateway pods: local-only vs global vs local front + global.
Run: python3 distributed_limits.py (stdlib only, under 1 second)
SharedStore stands in for Redis plus the Lua token bucket from the redis-lua page:
one lock makes check-and-debit atomic, exactly like one EVALSHA call.
"""
import random
import threading
import time
class TokenBucket:
"""Plain in-process token bucket. Not shared: every pod has its own copy."""
def __init__(self, rate: float, capacity: float) -> None:
self.rate, self.capacity = rate, capacity
self.tokens, self.ts = capacity, time.monotonic() # monotonic: no wall-clock jumps
def try_take(self, cost: float = 1.0) -> bool:
now = time.monotonic()
self.tokens = min(self.capacity, self.tokens + (now - self.ts) * self.rate)
self.ts = now
if self.tokens >= cost:
self.tokens -= cost
return True
return False
class SharedStore:
"""Fake Redis: one bucket per key, refill + debit under one lock (the Lua script)."""
def __init__(self, rate: float, capacity: float) -> None:
self.rate, self.capacity = rate, capacity
self.buckets: dict = {}
self.lock = threading.Lock()
self.calls = 0 # network round trips we would pay to Redis
self.down = False # flip to simulate a Redis outage
def take(self, key: str, cost: float = 1.0) -> bool:
if self.down:
raise TimeoutError("redis timeout")
with self.lock: # atomic like EVALSHA: nothing interleaves
self.calls += 1
b = self.buckets.setdefault(key, TokenBucket(self.rate, self.capacity))
return b.try_take(cost)
class GatewayPod:
"""Step 2-6 of Diagram 1. mode is 'local', 'global' or 'hybrid'."""
def __init__(self, mode: str, store: SharedStore, limit: float, pods: int) -> None:
self.mode, self.store = mode, store
if mode == "local":
# Wrong pick from Diagram 2: every pod enforces the full limit on its own.
self.local = TokenBucket(rate=limit, capacity=limit)
else:
# Hybrid front: a pod's fair share with 50% headroom for uneven load balancing.
share = limit / pods * 1.5
self.local = TokenBucket(rate=share, capacity=share)
# Fallback used only when Redis is down: strict per-pod share, so the total stays near the limit.
self.fallback = TokenBucket(rate=limit / pods, capacity=limit / pods)
self.fallbacks = 0
def handle(self, key: str) -> int:
if self.mode in ("local", "hybrid") and not self.local.try_take():
return 429 # Step 4a: fast reject, no Redis call
if self.mode == "local":
return 200
try:
ok = self.store.take(key) # Step 4b: atomic debit on shared Redis
except TimeoutError:
self.fallbacks += 1 # Failure path: degrade to local share and alert
return 200 if self.fallback.try_take() else 429
return 200 if ok else 429 # Step 5 -> 6a forward, 6b 429 + Retry-After
def run(mode: str, pods: int = 3, limit: float = 100, requests: int = 600, redis_down: bool = False):
store = SharedStore(rate=limit, capacity=limit)
store.down = redis_down
fleet = [GatewayPod(mode, store, limit, pods) for _ in range(pods)]
rng = random.Random(7)
codes = [rng.choice(fleet).handle("tenant:42") for _ in range(requests)] # LB spreads one key
return codes.count(200), store.calls, sum(p.fallbacks for p in fleet)
if __name__ == "__main__":
print("limit 100 per s for key tenant:42, 3 pods, 600 requests in one burst")
for mode in ("local", "global", "hybrid"):
ok, calls, _ = run(mode)
print(f"{mode:7s} admitted={ok:3d} redis_calls={calls}")
ok, calls, fb = run("hybrid", redis_down=True)
print(f"hybrid with Redis down: admitted={ok} redis_calls={calls} fallbacks={fb}")
ok6, _, _ = run("local", pods=6)
print(f"local with 6 pods: admitted={ok6} (limit multiplies with pod count)")Sample run on the box: local admitted 300 of 600 with 0 Redis calls (3 pods x 100); global admitted 100 with 600 Redis calls; hybrid admitted 100 with 150 Redis calls; hybrid with Redis down admitted 99 via the per-pod fallback; local with 6 pods admitted 577.
Working TypeScript
Run npx tsx distributed_limits.ts (no dependencies; the store is simulated with 2 ms latency). The fake store does refill and debit in one step, standing in for one EVALSHA call.
/**
* Distributed rate limits: local token bucket in front of a shared (Redis-like) store.
* Run: npx tsx distributed_limits.ts (no deps; the store is simulated with 2 ms latency)
* The fake store does refill + debit in one step, standing in for one EVALSHA Lua call.
*/
class TokenBucket {
private tokens: number;
private ts = performance.now(); // monotonic clock, safe against wall-clock jumps
constructor(private rate: number, private capacity: number) {
this.tokens = capacity;
}
tryTake(cost = 1): boolean {
const now = performance.now();
this.tokens = Math.min(this.capacity, this.tokens + ((now - this.ts) / 1000) * this.rate);
this.ts = now;
if (this.tokens >= cost) {
this.tokens -= cost;
return true;
}
return false;
}
}
const sleep = (ms: number) => new Promise<void>((r) => setTimeout(r, ms));
class FakeRedis {
calls = 0;
down = false; // when true, calls hang until the client timeout fires
private buckets = new Map<string, TokenBucket>();
constructor(private rate: number, private capacity: number) {}
async take(key: string): Promise<boolean> {
if (this.down) return new Promise<boolean>(() => {}); // never resolves, like a dead socket
await sleep(2); // network hop
this.calls++;
let b = this.buckets.get(key);
if (!b) this.buckets.set(key, (b = new TokenBucket(this.rate, this.capacity)));
return b.tryTake(); // single-threaded event loop: this block is atomic, like Lua
}
}
function withTimeout<T>(p: Promise<T>, ms: number): Promise<T> {
return Promise.race([p, sleep(ms).then(() => Promise.reject(new Error("redis timeout")))]);
}
class GatewayPod {
private front: TokenBucket; // Step 2: per-pod share with headroom
private fallback: TokenBucket; // strict share, used only while Redis is down
fallbacks = 0;
constructor(private redis: FakeRedis, limit: number, pods: number) {
this.front = new TokenBucket((limit / pods) * 1.5, (limit / pods) * 1.5);
this.fallback = new TokenBucket(limit / pods, limit / pods);
}
async handle(key: string): Promise<number> {
if (!this.front.tryTake()) return 429; // Step 4a: fast 429, no Redis call
try {
const ok = await withTimeout(this.redis.take(key), 20); // Step 4b: global debit
return ok ? 200 : 429; // Step 5 -> 6a forward or 6b 429 with Retry-After
} catch {
this.fallbacks++; // Failure path: local limit only, page someone
return this.fallback.tryTake() ? 200 : 429;
}
}
}
async function run(redisDown: boolean) {
const limit = 100, pods = 3;
const redis = new FakeRedis(limit, limit);
redis.down = redisDown;
const fleet = Array.from({ length: pods }, () => new GatewayPod(redis, limit, pods));
const codes = await Promise.all(
Array.from({ length: 600 }, (_, i) => fleet[i % pods]!.handle("tenant:42")),
);
const ok = codes.filter((c) => c === 200).length;
const fb = fleet.reduce((s, p) => s + p.fallbacks, 0);
console.log(`redisDown=${redisDown} admitted=${ok} redisCalls=${redis.calls} fallbacks=${fb}`);
}
await run(false);
await run(true);Sample run on the box: redisDown=false admitted 100, redisCalls=150, fallbacks=0. redisDown=true admitted 99, redisCalls=0, fallbacks=150 (each call gave up after the 20 ms timeout). tsc --strict: OK.
Interview Q&A
Your API runs on 20 gateway pods and the limit is 100 req/s per API key. Why not just set 100 on each pod?
Answer
Each pod only sees part of the traffic, so the client gets up to 20 x 100 = 2000 req/s, and the number changes whenever the fleet scales. Either share one counter (Redis) or divide the limit by the pod count, and the division is only right when the balancer spreads traffic evenly.
Why does the hybrid design still need Redis if the local bucket already rejects?
Answer
The local bucket only knows this pod. It rejects clear overage cheaply, but the sum of local buckets can still exceed the real limit, so the global debit is the source of truth. The local front just keeps abusive bursts from turning into Redis load.
Redis is down. Fail open or fail closed?
Answer
Default to fail-open to a strict local share (limit/N) with a short timeout and a circuit breaker, and alert. That keeps the API up with roughly the right total. Fail closed only where an overshoot is worse than an outage, such as payouts or SMS sending that costs money.
How do you avoid a Redis round trip on every request at very high RPS?
Answer
Lease tokens in batches: a pod takes 10 or 100 tokens from the global bucket in one call and spends them locally. Overshoot is bounded by pods x batch size. Alternatively, sync local counters to Redis every 100 ms and accept that window of error.
How do you rate limit across three regions?
Answer
Keep the hard limit per region in a regional Redis so the hot path stays local, then run an asynchronous global quota service that reallocates each region's share every few seconds from reported usage. You accept a small overshoot during each sync interval instead of paying WAN latency per request.
One tenant is so busy that its Redis key is hot. What do you do?
Answer
Shard the key (for example tenant:42 and a small shard id) and give each pod a shard, then sum or allocate per shard. Or use the local lease pattern so most requests never touch Redis. Also check that the hash tag is per identity, not one constant, or every key lands on the same Cluster slot.
How would you test that the distributed limit is actually correct?
Answer
Run a load test with many pods and a single key, then compare admitted count per window to the configured limit. Repeat after scaling the fleet up and down, with the balancer skewed, and with Redis latency and outages injected. Track admitted versus limit as a metric in production too.
Pros and cons
| Approach | Pros | Cons |
|---|---|---|
| Local per-pod limiter | No network hop. Keeps working when the store is down. Good for coarse abuse protection. | Effective limit is N x limit, or it depends on even balancing. Changes with autoscaling. Useless for paid quotas. |
| Global shared store (Redis + Lua) | Exact per-key limit across the fleet. One place to inspect and reset. | One round trip per request. Hot keys. Store outage needs a fallback policy. Cross-region latency if shared globally. |
| Local front plus global store | Exact limit with most rejects served locally. Degrades to local share on store failure. | Two configs to keep consistent. More moving parts to test. Small overshoot if you lease tokens in batches. |
Pitfalls
Draw 20 pods, a CDN, and one Redis. Mark where a scrape, a paying API key, and a GraphQL query get limited. Compute the effective limit if each pod allows 100 req/s. Then say what /login and GET /healthz do when Redis times out.
Go Deeper
- Stripe: Scaling your API with rate limiters
- Cloudflare: Counting things, a lot of different things
- Envoy global rate limiting (architecture overview)
- Envoy local rate limit filter
- envoyproxy/ratelimit, the reference global rate limit service
- Redis: Scripting with Lua
- Figma: An alternative approach to rate limiting
Related
- Rate Limiting: Token Bucket, Leaky Bucket & Sliding Window (
rate-limiting) - Token Bucket vs Leaky Bucket vs Sliding Window (
token-leaky-sliding-window) - Redis + Lua Atomic Rate Limiters (
redis-lua-atomic-rate-limiters) - HTTP 429, RateLimit Headers & Retry-After (
http-429-ratelimit-headers) - Fairness, Quotas & Noisy Neighbors (
fairness-quotas-noisy-neighbors)
Prev / Next
- Prev: Redis + Lua Atomic Rate Limiters (
redis-lua-atomic-rate-limiters) - Next: HTTP 429, RateLimit Headers & Retry-After (
http-429-ratelimit-headers) - Hub: Rate Limiting: Token Bucket, Leaky Bucket & Sliding Window (
rate-limiting)