Skip to main content

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​

PropertyTargetConsequence
Added latency< 2ms P99One in-region Redis round trip, or a local decision
AccuracyWithin ~1% of the configured limitApproximate algorithms are acceptable; unbounded overshoot is not
AvailabilityLimiter failure must not fail the APIFail-open with a local fallback
Scale1M req/s × ~4 rules = 4M checks/secCannot 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​

ComponentConcrete choiceWhy this one
Enforcement pointEnvoy ratelimit filter, or a gateway pluginEnforce at the edge, before the request consumes a thread in your service. Rejecting cheaply is the point
Counter storeRedis Cluster (or Aerospike for extreme scale)Single-threaded per shard, so Lua scripts are atomic without locks
AtomicityRedis Lua scripts (EVALSHA)Read-modify-write in one server-side operation; a pipeline is not atomic
Local fallbackGuava RateLimiter / golang.org/x/time/ratePreserves partial protection when Redis is unreachable
Rule distributionetcd or S3 + versioned poll, cached in-processRules change often; a Redis round-trip per rule lookup doubles latency
Cross-regionKafka + a per-key aggregator, or CRDT countersAsync convergence, no synchronous WAN hop

4. Algorithm selection (get here fast, then move on)​

AlgorithmMemory/keyBurst behaviourAccuracyVerdict
Fixed window counter8 BAllows 2× at a window boundaryPoorOnly for coarse quotas
Sliding window log8 B × N requestsExactPerfectToo much memory at scale
Sliding window counter16 BSmooth~1% errorGood default for quotas
Token bucket16 BConfigurable burstExactBest default for APIs
Leaky bucket (queue)Queue depthShapes to a constant rateExactUse 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.

  1. Shard the key itself. For keys above a detected threshold, split into key#0..key#N with each holding limit/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.
  2. 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 × instances of overshoot, bounded and known.
  3. 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​

EventFirst failureMitigation
Traffic doublesRedis shard CPU on Lua evalsGateway-side batching; more shards
One huge tenantSingle hot shard, all co-tenants degradeKey sharding with #N suffixes; dedicated buckets
Redis failover (~10s)All checks time out5ms timeout + local fallback bucket; alarm on ratelimit.degraded
Clock skew between gatewaysBuckets refill early or lateNTP/chrony discipline; pass now from a single source per region; tolerate ±1s
Mass 429 eventSynchronised client retry waveRetry-After from the script + decorrelated jitter in SDKs
Rule misconfigurationLegitimate traffic blocked globallyStage 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-After computed 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."