Distributed Rate Limiter
Sharpened prompt. Design a rate limiter for an API gateway handling 1M requests/sec across 5 regions, enforcing per-user, per-API-key, per-IP and per-endpoint rules, adding under 2ms to P99, and behaving sensibly when its own storage is unavailable.
The algorithm is the easy part and candidates spend too long on it. The hard parts are atomicity, geo-distribution, what happens when the limiter itself fails, and the tenant whose key is hot enough to melt a shard.
1. Problem framing
Functional requirements
- Multiple limit dimensions evaluated per request: user, API key, IP, endpoint, tenant.
- Rules configurable at runtime without redeploy.
- Standard client feedback:
429,Retry-After,X-RateLimit-*. - Burst tolerance — a client that has been quiet for a minute should be allowed a short spike.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Added latency | < 2ms P99 | One in-region Redis round trip, or a local decision |
| Accuracy | Within ~1% of the configured limit | Approximate algorithms are acceptable; unbounded overshoot is not |
| Availability | Limiter failure must not fail the API | Fail-open with a local fallback |
| Scale | 1M req/s × ~4 rules = 4M checks/sec | Cannot be one Redis instance |
Back-of-the-envelope
Checks: 1M req/s × 4 dimensions = 4M checks/sec
Redis: ~100k Lua evals/sec/shard -> ~40 shards minimum, 60 with headroom
State: 100M active keys × ~100 B = 10 GB (fits comfortably)
Cross-region: 5 regions, 80–150ms RTT — synchronous global counting is off the table
at a 2ms budget. This single number decides the architecture.
That last line is the whole design. Write it on the board early.
2. High-level architecture
The two loops matter: a fast synchronous loop entirely inside one region, and a slow asynchronous loop that reconciles regional usage globally. Nothing crosses a region boundary on the request path.
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Enforcement point | Envoy ratelimit filter, or a gateway plugin | Enforce at the edge, before the request consumes a thread in your service. Rejecting cheaply is the point |
| Counter store | Redis Cluster (or Aerospike for extreme scale) | Single-threaded per shard, so Lua scripts are atomic without locks |
| Atomicity | Redis Lua scripts (EVALSHA) | Read-modify-write in one server-side operation; a pipeline is not atomic |
| Local fallback | Guava RateLimiter / golang.org/x/time/rate | Preserves partial protection when Redis is unreachable |
| Rule distribution | etcd or S3 + versioned poll, cached in-process | Rules change often; a Redis round-trip per rule lookup doubles latency |
| Cross-region | Kafka + a per-key aggregator, or CRDT counters | Async convergence, no synchronous WAN hop |
4. Algorithm selection (get here fast, then move on)
| Algorithm | Memory/key | Burst behaviour | Accuracy | Verdict |
|---|---|---|---|---|
| Fixed window counter | 8 B | Allows 2× at a window boundary | Poor | Only for coarse quotas |
| Sliding window log | 8 B × N requests | Exact | Perfect | Too much memory at scale |
| Sliding window counter | 16 B | Smooth | ~1% error | Good default for quotas |
| Token bucket | 16 B | Configurable burst | Exact | Best default for APIs |
| Leaky bucket (queue) | Queue depth | Shapes to a constant rate | Exact | Use when the upstream needs a smooth rate, not just a bounded one |
The fixed-window boundary bug is worth stating concretely because it is the classic gotcha: with a 100/minute limit, a client sends 100 requests at 11:00:59 and 100 more at 11:01:00 — 200 requests in one second, both windows technically respected.
Token bucket is the right default: it expresses "1000 requests/minute with a burst of 100" naturally, which is what API products actually sell.
5. The toughest parts
5.1 Atomicity: the check-then-set race
Why it's hard. The naive implementation is GET count; if count < limit: INCR. Between the GET and the INCR, a thousand other requests do the same thing. Every one of them reads a value below the limit and every one proceeds. Under concurrency, the limit is not merely exceeded — it is exceeded by an amount proportional to your concurrency, which is exactly when you needed it most. MULTI/EXEC does not fix this either: it batches, but the decision logic still lives in the client.
Solution — push the whole decision into a Lua script. Redis executes a script atomically on a single thread; nothing interleaves.
-- token_bucket.lua
-- KEYS[1] = bucket key
-- ARGV[1] = capacity, ARGV[2] = refill_per_sec, ARGV[3] = now_ms, ARGV[4] = cost
local capacity = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local cost = tonumber(ARGV[4])
local b = redis.call('HMGET', KEYS[1], 'tokens', 'ts')
local tokens = tonumber(b[1])
local ts = tonumber(b[2])
if tokens == nil then tokens = capacity; ts = now end
-- Lazy refill: no background timers, no cron. Compute what would have accrued.
local elapsed = math.max(0, now - ts) / 1000.0
tokens = math.min(capacity, tokens + elapsed * rate)
local allowed = 0
if tokens >= cost then
tokens = tokens - cost
allowed = 1
end
redis.call('HMSET', KEYS[1], 'tokens', tokens, 'ts', now)
-- TTL = time to refill from empty, so idle keys evict themselves
redis.call('PEXPIRE', KEYS[1], math.ceil(capacity / rate * 1000) + 1000)
local retry_after = 0
if allowed == 0 then retry_after = math.ceil((cost - tokens) / rate) end
return { allowed, math.floor(tokens), retry_after }
Two details that signal experience. Lazy refill: never run a background job to add tokens to millions of buckets — compute the accrual from the elapsed time on read. Self-expiring keys: the PEXPIRE means an idle client's bucket disappears on its own, which is what keeps 100M keys from becoming 10B.
Pass now from the caller rather than using redis.call('TIME') so the script stays deterministic and replica-safe.
5.2 Cross-region enforcement without a 150ms WAN hop
Why it's hard. "1000 requests/minute per API key" is a global promise, but the key's traffic arrives in 5 regions at once. A single global Redis means 80–150ms added latency on every request — 50× the budget. Independent per-region limits mean a client that spreads traffic across regions gets 5× the limit.
Solution — local enforcement, async global reconciliation. Each region enforces against a local quota share and reports consumption asynchronously; a global aggregator redistributes shares based on where the traffic actually is.
# Every 250ms, each region publishes what it consumed and pulls its next share.
async def reconcile(region: str):
while True:
used = await redis.eval(DRAIN_DELTAS, keys=[f"deltas:{region}"])
await kafka.send("quota.usage", {"region": region, "used": used, "t": now_ms()})
# The aggregator computes each region's share of the global limit,
# weighted by recent demand: a region taking 70% of traffic gets ~70% of quota.
shares = await kafka_compacted_topic.get("quota.shares")
await redis.eval(APPLY_SHARES, keys=["quota:local"], args=[shares[region]])
await asyncio.sleep(0.25)
Be explicit about the error this admits: with a 250ms reconciliation interval, a client can briefly exceed the global limit by up to regions × per_region_burst. Then justify it — "for a 1000/minute quota, a transient overshoot of a few dozen requests is commercially irrelevant; 150ms of added latency on every API call is not." That sentence is the answer the interviewer is listening for.
When approximation is not acceptable — a hard regulatory cap, or a paid-per-call downstream — pin the key to a home region (hash the API key to a region) and pay the WAN cost only for that tenant's out-of-region traffic. Naming the escape hatch is as important as the default.
CRDT alternative: model each counter as a G-Counter (a per-region vector, global value = sum). Regions merge by taking the element-wise max. It converges without coordination and never loses an increment, but it can only count up, so it fits fixed windows, not token buckets that refill.
5.3 Fail-open or fail-closed?
Why it's hard. Redis becomes unreachable. Block everything, and a rate limiter — a protective component — has caused a total outage. Allow everything, and an ongoing attack proceeds unthrottled, potentially taking down the services the limiter exists to protect. Both answers are defensible, which is precisely why "it depends" is not an answer. Say what it depends on.
Solution — a per-rule policy, defaulting to fail-open with degraded local enforcement.
func (l *Limiter) Allow(ctx context.Context, key string, rule Rule) Decision {
ctx, cancel := context.WithTimeout(ctx, 5*time.Millisecond) // hard budget
defer cancel()
res, err := l.redis.Allow(ctx, key, rule)
if err == nil {
l.local.Observe(key, res) // keep the local bucket roughly warm
return res
}
l.metrics.Inc("ratelimit.degraded") // alarm on this, always
switch rule.FailureMode {
case FailClosed: // auth endpoints, payment mutations
return Deny{Reason: "limiter_unavailable", Retry: 1 * time.Second}
default: // everything else
// Local in-process bucket, sized limit/instances × 1.2 for safety.
return l.local.Allow(key, rule.PerInstanceShare())
}
}
The five-millisecond timeout is doing as much work as the policy. A Redis that is slow rather than dead is the more dangerous failure: without a tight deadline, request threads pile up waiting on the limiter and the gateway falls over from the thing meant to protect it.
The rule of thumb to state: fail open when the limiter protects capacity (throughput throttling), fail closed when it protects correctness or money (login attempts, OTP sends, payment retries, expensive third-party calls).
5.4 The hot tenant that melts one shard
Why it's hard. Rate limit keys shard by the limited entity. One enterprise customer sends 300k req/s; their key lives on exactly one Redis shard, which now handles 300k Lua evals/sec against a ~100k ceiling. That shard's latency spikes, your 5ms timeout starts firing, and every tenant on that shard degrades. The rate limiter has become the noisy-neighbour problem it exists to prevent.
Solution.
- Shard the key itself. For keys above a detected threshold, split into
key#0..key#Nwith each holdinglimit/N, and have the caller hash by client instance ID (not randomly — random selection means one instance can be denied while another has tokens). Load spreads across N shards. - Batch at the gateway. Instead of one Redis call per request, each gateway instance acquires tokens in blocks of 50 and serves 50 requests locally. Redis traffic drops 50×. The cost is up to
50 × instancesof overshoot, bounded and known. - Two-tier for the biggest tenants. Above a certain size, a customer gets a dedicated in-process bucket per gateway with a much slower reconciliation loop — the same idea as 5.2, applied per tenant.
Batching is the highest-leverage move and generalises: the rate limiter is itself a system that needs rate limiting, and the answer is amortisation.
5.5 Evaluating many rules without multiplying latency
Why it's hard. A single request may match four rules (global, per-key, per-IP, per-endpoint), plus tenant overrides and endpoint-specific costs. Four sequential Redis round trips is 4× the latency budget. Worse, if rule 1 allows and rule 4 denies, you have already consumed tokens from three buckets for a request that never runs — a slow leak that under-serves clients over time.
Solution — evaluate all dimensions in one atomic script, and only commit if all pass.
-- multi_limit.lua: check every bucket first, deduct only if all allow.
local n = #KEYS
local states = {}
for i = 1, n do
local allowed, tokens = check_bucket(KEYS[i], ARGV[i]) -- pure, no writes
if allowed == 0 then
return { 0, i } -- deny, report which rule tripped
end
states[i] = tokens
end
for i = 1, n do
commit_bucket(KEYS[i], states[i]) -- all-or-nothing deduction
end
return { 1, 0 }
This requires all four keys to live on the same Redis slot, which you force with a hash tag: {user:42}:global, {user:42}:endpoint:search. Everything in braces determines the slot, so the whole evaluation is one shard, one script, one round trip. Redis hash tags are a small detail that demonstrates real operational familiarity.
Rules themselves live in an in-process cache refreshed from etcd every few seconds — never fetch configuration on the request path.
5.6 Telling the client the truth (and avoiding the retry stampede)
Why it's hard. A bare 429 with no guidance produces the worst possible client behaviour: immediate retry. A thousand throttled clients retrying in a tight loop generate more load than the traffic you rejected, and because they all got throttled at the same instant, they retry in a synchronised wave — forever.
Solution — precise headers plus mandatory jitter.
HTTP/1.1 429 Too Many Requests
RateLimit-Limit: 1000
RateLimit-Remaining: 0
RateLimit-Reset: 27 ; seconds until tokens are available (RFC 9239 draft)
Retry-After: 27
X-RateLimit-Scope: user
Return Retry-After from the Lua script (it already computed exactly when enough tokens accrue — see 5.1), not from a constant. And in your official SDKs, implement decorrelated jittered backoff, which is what actually prevents the synchronised wave:
# AWS's decorrelated jitter: no two clients converge on the same retry instant.
sleep = min(cap, random.uniform(base, sleep * 3))
Note also that returning RateLimit-Remaining on successful responses lets well-behaved clients self-throttle before they ever hit a 429 — the cheapest request to serve is the one the client decided not to send.
5.7 Distinguishing throttling from abuse
Why it's hard. A rate limiter answers "is this client over their quota?" It cannot answer "is this a botnet rotating through 50,000 residential IPs, each staying just under the per-IP limit?" Per-key limits are structurally blind to distributed low-and-slow attacks, and tightening the limit to catch them punishes legitimate users.
Solution — layer a detector above the limiter. Keep the deterministic limiter fast and dumb; run anomaly detection asynchronously on the request stream (Flink or a sketch-based detector) looking at aggregate signals: request entropy, ASN concentration, TLS fingerprint (JA4) clustering, and the ratio of new to returning identities. When a cluster is flagged, push a dynamic rule through the control plane — a tighter limit scoped to that ASN or fingerprint — which the gateway picks up within seconds.
The framing to use: "Rate limiting is a quota mechanism, not a security control. It bounds cost. Abuse detection is a separate, slower, stateful system that expresses its conclusions as rate limit rules." That separation is the senior answer.
6. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
| Traffic doubles | Redis shard CPU on Lua evals | Gateway-side batching; more shards |
| One huge tenant | Single hot shard, all co-tenants degrade | Key sharding with #N suffixes; dedicated buckets |
| Redis failover (~10s) | All checks time out | 5ms timeout + local fallback bucket; alarm on ratelimit.degraded |
| Clock skew between gateways | Buckets refill early or late | NTP/chrony discipline; pass now from a single source per region; tolerate ±1s |
| Mass 429 event | Synchronised client retry wave | Retry-After from the script + decorrelated jitter in SDKs |
| Rule misconfiguration | Legitimate traffic blocked globally | Stage rules in shadow mode (log-only) before enforcing; one-click rollback |
Shadow mode deserves the emphasis: every new rule runs for a day recording what it would have blocked, and only then flips to enforcing. It is the single practice that keeps rate limiters from causing outages, and almost nobody mentions it.
7. Cheat sheet
- Algorithm: token bucket with lazy refill. Sliding-window counter for quota-style limits. Never fixed window.
- Atomicity: one Lua script per decision. Pipelines are not atomic.
- Multi-rule: hash tags force co-location; check-all-then-commit-all in one script.
- Cross-region: enforce locally, reconcile via Kafka every 250ms; pin to a home region only when the cap is hard.
- Failure: 5ms timeout, fail-open with a local bucket by default, fail-closed for auth and money.
- Hot keys: batch token acquisition at the gateway — 50× fewer Redis calls for bounded overshoot.
- Clients:
Retry-Aftercomputed from actual refill time, plus decorrelated jitter in the SDK. - The one-liner: "Enforcement is local and synchronous, correctness is global and asynchronous, and the interesting engineering is entirely in the gap between them."