Design a Distributed Rate Limiter
DifficultyMedium | HelloInterview: problem breakdown
Problem Statement
Build a rate limiter for a large API platform. Clients should be limited to N requests per time window. The system runs across multiple data centers and hundreds of API gateway nodes.
๐Real-world & security angle: The rate limiter is where systems design and security overlap most directly โ it's the front-line defense against brute-force, credential stuffing, scraping, DoS, and cost-abuse (the LLM "denial-of-wallet" problem). The classic algorithm family is worth knowing cold: token bucket (allows bursts, refills steadily โ what most APIs and AWS use), leaky bucket (smooths to a constant rate), fixed window (simple but has a boundary-spike flaw โ 2N requests across a window edge), and sliding window (fixes the edge problem). The hard distributed twist: with hundreds of gateway nodes you can't keep per-node counters, so you centralize state in Redis (often with a Lua script for atomic check-and-increment) โ and then face the real-world decision interviewers probe: fail-open or fail-closed when Redis is down? Most APIs fail open (allow traffic) to preserve availability, accepting that the limiter briefly stops protecting โ a genuine security-vs-availability tradeoff. GitHub, Stripe, and Cloudflare all publish how they do this.
Requirements
Functional
- Enforce per-client (API key / user ID) request limits
- Configurable:
N requests per window(per second, minute, hour, day) - Return
429 Too Many RequestswithRetry-Afterheader when exceeded - Different limits per endpoint tier (free vs paid)
- Soft vs hard limits (burst allowance)
Non-Functional
- Decision latency: < 5ms (inline with every API request)
- Very high availability โ rate limiter failure should fail open (don't block all traffic)
- Accuracy: allow N ยฑ small epsilon (exact enforcement is less important than low latency)
- Distributed: decision must be consistent across all gateway nodes for the same client
Scale
- 500k API requests/s
- 10M active API keys
- Distributed across 3 regions
Algorithms
Token Bucket
bucket(capacity=100, refill_rate=10/sec)
on request:
if bucket.tokens > 0: bucket.tokens -= 1; ALLOW
else: DENY (429)
background: refill tokens at rate until capacityallows bursting (up to capacity), smooth average rate
state per client, slightly complex to implement distributed
Sliding Window Log
- Store timestamp of each request for the client
- On new request: count timestamps within the last window โ compare to limit
- Pros: very accurate
- Cons: high memory (store every timestamp), slower for high-rate clients
Sliding Window Counter (recommended for distributed)
count = current_window_count + (previous_window_count ร overlap_fraction)- Two counters (current and previous 1-min window) + fractional interpolation
- Very low memory (2 counters per client), fast, accurate enough
- Memory: 10M clients ร 2 counters ร 8 bytes = 160 MB โ fits in Redis
Fixed Window Counter
- Simplest:
INCRkey{client}:{minute}in Redis; compare to limit - Problem: boundary spike โ a client can make 2N requests in 2 seconds straddling a window boundary
- Acceptable for loose limits; not for strict enforcement
Core Design
Client โ API Gateway โ Rate Limit Middleware
โ
Redis Cluster (counters)
โ
(local in-memory cache, sync every 100ms)Distributed Counter with Redis
def is_allowed(client_id: str, limit: int, window_sec: int) -> bool:
key = f"rl:{client_id}:{int(time.time() // window_sec)}"
pipe = redis.pipeline()
pipe.incr(key)
pipe.expire(key, window_sec * 2)
count, _ = pipe.execute()
return count <= limit- Lua script ensures atomicity: increment + check in single operation
- Redis single-threaded: no race conditions
Multi-Region Challenge
Sharing a single Redis across 3 regions introduces 50-80ms cross-region latency โ unacceptable for a <5ms decision.
Solutions
each region has its own Redis; counters are only eventually consistent. Acceptable for most use cases โ a client could exceed limits by up to N ร num_regions briefly.
each API gateway node keeps a local cache of remaining tokens, syncs to Redis every 100ms. Very fast (in-process), slightly over-counts at boundary.
use incr locally in-process, batch-sync to Redis every N requests or T milliseconds. Efficient at scale.
Fail-Open Design
If Redis is unavailable:
- Default to ALLOW (fail open) โ better to let some excess traffic through than to block all legitimate traffic
- Set a circuit breaker: if Redis latency > 5ms, bypass rate limiting
- Alert on Redis unavailability
Key Design Decisions
| Decision | Choice | Reason |
|---|---|---|
| Algorithm | Sliding window counter | Low memory, fast, accurate enough |
| Storage | Redis (in-memory) | Sub-millisecond latency, atomic INCR |
| Granularity | Per-client + per-endpoint | Allows tiered limits |
| Failure mode | Fail open | Availability > perfect enforcement |
| Multi-region | Local Redis + async sync | Latency wins over exactness |
Rate Limit Headers (RFC 6585)
X-RateLimit-Limit: 1000
X-RateLimit-Remaining: 450
X-RateLimit-Reset: 1709812800
Retry-After: 23Security Considerations
| Threat | Mitigation |
|---|---|
| Rate limit bypass by rotating IPs | Limit on API key/user ID, not just IP |
| Credential stuffing | IP-based limits AND user-based limits (defense in depth) |
| DDoS against rate limiter itself | Rate limiter must be in-process or very close (not a remote call) |
Fake X-Forwarded-For header to spoof IP | Trust only the leftmost IP from known trusted load balancers; validate header chain |
| Enumeration via Retry-After | Slightly randomize Retry-After to prevent attackers from timing their retry exactly |
| Key exhaustion (very high cardinality) | TTL on all Redis keys; eviction policy allkeys-lru for Redis |
Security use case for rate limiting
- Brute-force login: 5 attempts per 15min per IP + per account
- Password reset: 3 requests per hour per user
- SMS OTP: 5 per day per phone number
- API scraping: graduated limits; escalate to CAPTCHA then block
Interview Tips
- Explain the algorithm trade-offs before picking one โ interviewers want to see you reason, not just recite.
- The distributed consistency problem is the key challenge here โ how do you make a decision in <5ms when your state is distributed?
- Fail-open is a deliberate choice โ security engineers might disagree (fail-closed for security), but for a general API rate limiter, availability wins. State this trade-off explicitly.
- Show you know the RFC 6585 response headers โ it's a detail that signals production experience.