Distributed Rate Limiter
Worked solution

Distributed Rate Limiter — a worked solution

Enforce a per-key request limit across a fleet of enforcers — accurately, in under a millisecond, without becoming the outage.

Try it yourself first.

You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.

Open the Distributed Rate Limiter workspace

The problem

Build a distributed rate limiter: cap how many requests a given key — an API key, a user id, an IP — may make in a window, and enforce that cap across a whole fleet of enforcer instances, on the synchronous path of every request, in well under a millisecond, without ever becoming the outage you were trying to prevent.

The one-sentence version ("count requests in Redis, reject over the limit") is right and useless. The entire problem is the gap between a global invariant ("key X ≤ R req/s") and the local vantage points that must enforce it — 200 enforcer pods, each seeing only a slice of X's traffic. Close that gap naively and a client fans out across your fleet and gets 200× the limit; close it with a single shared counter and you've put a network round trip (and a single point of failure) on every request. This problem is a distributed-consistency problem wearing a counter's clothes.

We hold everything else fixed and go deep on that core: the token-bucket algorithm as an atomic Redis Lua script, a central-authoritative counter with a per-instance fail-open fallback, and the four tensions that define the space — accuracy vs latency vs availability, and what you do when the store dies.

The reference architecture

Reference architecture for Distributed Rate Limiter: 8 components — Client, Edge LB, API Gateway · Enforcer, Service · Rate Limiter, KV Store · Counter (Redis Cluster), Protected API · Upstream, Control Plane · Rate Rules, Monitoring — connected by 10 flows.ClientclientEdge LBL7 LB, consistent-hash …API Gateway · EnforcerEnvoy (ratelimit filter…Service · Rate LimiterLyft-style ratelimit (G…KV Store · Counter (R…Redis 7 Cluster (16 sha…Protected API · Upstr…The service the limiter…Control Plane · Rate …xDS / config store, hot…MonitoringPrometheus + Alertmanag…
8 components, 10 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Client

Issues requests carrying an identity the limiter can key on — an API key (Authorization: Bearer k_...), a user id, or (for anonymous traffic) the source IP. A well-behaved client reads the rate-limit signal on every response (RateLimit: "default";r=…;t=… and, on a 429, Retry-After) and paces itself: it stops sending when r hits 0 and waits t/Retry-After seconds before retrying, with exponential backoff AND jitter so a fleet of clients doesn't re-synchronize into a thundering herd.

Why it exists. Every request enters here — there is no other ingress. The interesting design fact is that the CLIENT is half of the rate-limiter's correctness story: the server can return 429s all day, but if clients retry immediately and in lockstep, those retries *sustain* the overload the limiter exists to shed (Cloudflare's Nov-2023 DR recovery: failing calls formed a thundering herd until limits were imposed on the reconnect surge). Google SRE pushes this further — a well-built client runs *adaptive throttling* (p(reject)=max(0,(requests−2·accepts)/(requests+1))) and sheds its own traffic locally when it sees a high rejection fraction, because even *rejecting* a request costs the server CPU (RFC 6585 says so explicitly).

When it fails. The retry storm. A buggy client — the canonical case is Cloudflare's Sept-2025 dashboard outage, where a React useEffect with an unstable dependency re-fired an API call on every render — turns one client into thousands of requests/sec against a first-party service that had no protective limit. Mitigation is defense in depth: the limiter caps the blast radius per key, AND clients must backoff+jitter, AND the protected service self-protects regardless of quota.

Edge LBL7 LB, consistent-hash on rate-limit key

Terminates TLS and spreads inbound requests across the enforcer fleet. The load-bearing configuration choice is the balancing algorithm: instead of round-robin, it uses consistent-hash routing on the rate-limit key (API key / client IP), so all traffic for a given key lands on the *same* enforcer instance. In steady state (global store healthy) this doesn't change correctness — the global counter is authoritative regardless of routing — but it is exactly what makes the per-instance *fallback* bucket correct when the store is unreachable, and it pins a single attacker to a single enforcer instead of letting them fan out across all 200.

Why it exists. This is the same insight that makes Cloudflare's edge design work: because anycast steers one IP to one PoP, a PoP-local counter is a good proxy for a global one. We borrow it deliberately for the fallback path. Round-robin would be simpler, but then a store outage degrades to a per-replica bucket that an attacker fan-outs across (see rate-limiter:store-unreachable-failopen): 200 instances × R each = 200×R admitted. Consistent-hash pins the key so the local bucket sees the key's *whole* stream, making the degraded limit ≈ R, not 200×R.

When it fails. Routing churn (scale event, instance eviction, ring rebalance) moves a key to a new enforcer whose local bucket is empty — a brief over-admission during the reshuffle. Also, consistent-hash concentrates a celebrity key onto one enforcer; if that key is genuinely huge, that instance is hot regardless of the global store. Both are bounded and acceptable versus the alternative (round-robin's 200× fallback blowup).

API Gateway · EnforcerEnvoy (ratelimit filter + local_rate_limit)

The point where the limit is enforced. For each request it builds a descriptor — an ordered vector of key/value tuples like (api_key, k_abc)→(route, /v1/search)→(tier, free) — and calls the Rate Limiter Service over gRPC (ext_authz-style, ~1–2 ms, no retry, tight timeout). On OK it forwards the request to the upstream and stamps RateLimit/RateLimit-Policy headers on the response; on OVER_LIMIT it returns 429 immediately and never touches the upstream (with Retry-After + RateLimit;r=0). It ALSO runs a cheap in-process token bucket as a first tier: it sheds obvious floods at line rate (protecting the shared limiter from being overwhelmed) and, if the limiter/store is unreachable, it *fails open onto that local bucket* rather than 429ing everyone.

Why it exists. Enforcing at the gateway means the expensive request never reaches the upstream when it's over budget — the whole value proposition. Envoy is the reference because it ships *both* tiers of the canonical pattern: local_rate_limit (an in-process token bucket, per-instance and therefore approximate) in front of the global gRPC ratelimit service (Redis-backed, accurate across the fleet). Envoy's own docs recommend running both — 'local rate limiting can be used in conjunction with global rate limiting to reduce load on the global rate limit service.' We reject enforcing inside each upstream service (N re-implementations, and you can't trust a service that hasn't checked the caller) and reject a pure local limiter (per-instance → N× the limit, see fanoutOveradmit).

When it fails. The gateway is on the synchronous path of every request, so its own latency is pure overhead on p50 AND p99 — if the global check's tail (a hot Redis shard, a GC pause, a conn-pool stall) climbs, it shows up as tail latency on the *entire* API. Mitigation: local fast path absorbs bursts; 2 ms hard timeout on the check → fail open on timeout; circuit-break the limiter call at high error rate. Pages on-call when check_latency_ms p99 > 5 ms or fallback_active > 0 for more than 60 s.

Service · Rate LimiterLyft-style ratelimit (Go, stateless)

Stateless decision service. Receives a (domain, descriptor) from the gateway, resolves the matching rule (limit, window, algorithm, tier override) from its in-memory rule cache — most-specific-match over the descriptor tree, pushed by the Control Plane — then executes the token-bucket Lua script atomically against the Redis shard that owns the key. It returns {allowed, limit, remaining, retry_after_s, reset_s} in well under a millisecond of its own compute. All durable state lives in Redis and all policy lives in the rule cache, so the service itself holds nothing — it scales out horizontally and any replica can serve any descriptor.

Why it exists. Separating the decider from the enforcer (gateway) lets us evolve the algorithm, the rule schema, and the Redis topology without touching the data plane, and lets the decider be language-owned in one place (the Lua scripts, the descriptor matcher, the Redis client) rather than re-implemented in every gateway. It is stateless on purpose: the moment it held per-key state it would need its own replication and failover, duplicating Redis. The canonical open-source shape is Lyft's envoyproxy/ratelimit; we keep its descriptor model but run token-bucket Lua (burst-friendly, O(1)) instead of its default fixed-window counters (which carry the 2× boundary bug).

When it fails. The decision tier saturates under a traffic spike (rate-limiter:limiter-tier-overload): the gRPC check queues, and because the check is synchronous the latency lands on every request. Mitigation is the two-tier design — the gateway's local bucket sheds the flood before it reaches this service — plus autoscaling on check-queue depth and a strict gateway-side timeout that fails open rather than waiting. A subtler failure is a rule-cache that's silently stale (Control Plane push didn't land): the service enforces yesterday's limits; detected by a rule-version metric diffed against the Control Plane.

KV Store · Counter (Redis Cluster)Redis 7 Cluster (16 shards, RF=2, AOF everysec)

The authoritative distributed state: one token-bucket entry per key — {tokens: float, last_refill_ts} — packed into a single value so the Lua script touches one key atomically. Refill is lazy (computed from the timestamp delta on each check: tokens = min(capacity, tokens + elapsed × refill_rate)), so there is no background dripper. Keyspace is split across 16384 hash slots (CRC16(key) mod 16384) divided among the shards; hash tags {key} colocate a key's minute/hour buckets on one slot so a multi-window Lua script stays single-slot. Idle buckets carry a short TTL so cold keys evict; the whole active working set is tiny (~0.3 GB for 5M keys at 64 B).

Why it exists. A distributed rate limiter is *defined* by this component: it's the shared truth that lets 200 enforcers converge on one global count for a key. Redis is the industry default (Stripe, GitHub, Envoy/Lyft, Kong-redis) because it's single-threaded per shard — which is exactly what makes an atomic Lua read-modify-write cheap and correct — and because ~100k–200k un-pipelined atomic ops/sec per shard is plenty for a check-per-request workload once sharded. We reject Memcached (no atomic multi-step RMW — the reason Cloudflare's edge had to use an approximation and old GitHub had eviction bugs) and reject a SQL row-per-key (a read+write to Postgres per request is orders of magnitude too slow).

When it fails. Two signature failures. (1) Hot shard (rate-limiter:hot-key-shard): a single hot key is one slot on one core — adding nodes can't help, so that shard saturates while the other 15 idle; fix by sharding the counter or absorbing at the local tier. (2) Failover free-burst (rate-limiter:counter-store-failover): async replication means the promoted replica is missing the last ~1 s of increments, so the counter regresses and clients get a free burst during the ~12 s failover — the worst moment for the downstream. Pages on-call on redis_failover_total, per-shard ops_per_sec imbalance, and aof_last_write_status != ok.

Protected API · UpstreamThe service the limiter shields

The upstream the rate limiter exists to protect — the application API, or a fragile shared dependency behind it (a database, a third-party integration, an expensive search cluster). It only ever sees requests the gateway has already ADMITTED; an over-limit request is rejected at the gateway with a 429 and never arrives here. That is the entire point of enforcing at the edge rather than inside the upstream: the expensive work is never started for traffic that's over budget.

Why it exists. Modeled explicitly so the two decisive journeys are visible on the diagram: the *allow* path reaches this node (through the gateway's forward), and the *deny* path terminates at the counter store and this node stays cold. It also anchors the rate-limiting-vs-load-shedding distinction: the rate limiter protects this upstream from a single noisy tenant (per-key fairness), but even with every caller in-quota, 10k in-quota callers can spike together — so this upstream must *also* self-protect (shed by criticality, reject early with 503 when in-flight work exceeds a threshold) rather than trust the quota layer alone.

When it fails. Overload despite the limiter — because per-key limits don't bound *aggregate* in-quota traffic. This is why the design pairs the (external-facing, per-key) rate limiter with (internal) load shedding: the two are different mechanisms for different threats. If this upstream falls over, clients retry, and without client backoff+jitter the retries pin it — the same thundering-herd trap the whole design guards against.

Control Plane · Rate RulesxDS / config store, hot reload

The source of truth for *policy*: which descriptor gets which limit — per route, per plan tier (free 100/min vs premium 1000/min), per-key overrides for a specific abusive or whitelisted caller, and the algorithm/window for each. It distributes those rules to the Rate Limiter fleet (and the coarse local-limit + shed config to the gateways) via hot reload — a file/runtime watch with an atomic symlink-swap, or an xDS gRPC push of RateLimitConfig resources — so a limit change takes effect fleet-wide without a deploy or a restart. It also owns the safe-rollout primitives: shadow_mode (evaluate the rule, emit 'would-have-blocked' stats, but always return OK) and rule replaces for deterministic tier overrides.

Why it exists. Limits are operational config that changes constantly (a new tier, an emergency clamp on an abusive key, a route that suddenly needs protecting during an incident) — hard-coding them into the enforcers would mean a deploy per change, exactly when you can least afford one. Separating policy (this control plane) from mechanism (the limiter's Lua + Redis counters) is what lets an on-call engineer tighten a limit in seconds. We keep it a small, separate authority rather than folding rules into the limiter so that a rules push and a counter operation are independent failure domains.

When it fails. Two shapes. (1) The control plane is DOWN (rate-limiter:control-plane-down): you cannot push a new limit — during an attack you can't clamp the abusive key, and during a good deploy you can't roll out a new tier; enforcers keep running on stale-but-safe rules. (2) A BAD rule is pushed (limit=0, or a match that's too broad): without shadow_mode + canary it fail-closes a whole route fleet-wide in seconds — the config-push blast radius (Cloudflare's 2019 regex-CPU and Knight Capital's partial deploy are the cautionary class). Mitigation: shadow → canary → progressive rollout, and an instant rollback path.

MonitoringPrometheus + Alertmanager + Grafana

Pull-based metrics from the gateway, rate limiter, and Redis every 15 s: allowed_total / denied_total per descriptor, check_latency_ms p50/p99, fallback_active (local bucket in use because the store was unreachable), per-shard redis_ops_per_sec (hot-shard detection), redis_failover_total, aof_last_write_status, and rule_version (per enforcer, diffed against the control plane to catch stale-rule drift). Alertmanager fires PagerDuty on the specific signals each component's failureMode names — check p99 > 5 ms, fallback active > 60 s, shard ops imbalance, failover events, rule-version skew.

Why it exists. Every failure mode in this design names a page-on-call signal; without an alerting plane those pages have nowhere to fire. A rate limiter is a load-bearing dependency whose *own* degradation is invisible to users until it either leaks abuse (fail-open during a store outage) or blocks everyone (fail-closed) — so the tell-tale metrics (fallback active, denied-rate spike, check-latency tail) are the early-warning system. It's separate infrastructure because it must survive when the thing it's watching doesn't.

When it fails. Cardinality explosion if someone adds the raw rate-limit key (millions of them) as a metric label — capped per-metric with a series_per_target alarm. And the watcher-of-the-watchman problem: if Alertmanager is down, pages vanish silently — mitigated with an external dead-man's-switch heartbeat.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Surface the constraints up front:

  • What are we limiting, and per what key? Requests per second/minute, keyed on API key (the real "who", tied to billing) and/or IP (anonymous abuse, coarse and spoofable). Route and plan-tier are additional descriptor dimensions. Drives the descriptor + rule model.
  • One global limit, or per-region? Assume a single logical limit per key, enforced from many instances in one region. True cross-region global limits are a deliberate deep-dive (they can't sit on the synchronous path).
  • Which algorithm? Token bucket (burst-friendly, O(1), the default). Fixed / sliding-window are alternatives we justify against it.
  • Accuracy tolerance? A little over-admission is acceptable; blocking legitimate traffic is not. This is the load-bearing product call — it makes the whole thing AP, not CP.
  • Latency budget? The check is on every request's critical path. Target < 1 ms added at p50 and a bounded tail; it must never become the bottleneck it protects against.
  • What happens when the counter store is down? Fail open (allow, accept brief abuse) — chosen deliberately, with a local fallback so the floodgates don't open fully.
  • What does the client see on a limit hit? 429 + Retry-After + RateLimit headers; clients must back off with jitter.

Assumptions to state:

  • ~1,000,000 req/s peak across the fleet, ~300,000 avg; one atomic check per request.
  • ~5,000,000 distinct active keys; token-bucket state is ~64 B/key → the whole working set is ~0.3 GB (tiny — this is a throughput and latency problem, not a storage one).
  • 200 enforcer instances; ~16 Redis shards (~10 needed for throughput, headroom to 16).
  • ~2% of requests are over-limit and get 429'd; the other 98% are forwarded.

02Functional reqs

What must this system actually do?

  • Decide allow/deny for each request against the limit for its key — atomically, so concurrent requests can't both slip past the last token.
  • Enforce the decision at the gateway: allowed → forward to the upstream; denied → 429 without touching the upstream.
  • Report quota to the client on every response: RateLimit (remaining, reset) and, on denial, Retry-After.
  • Configure limits by descriptor — per route, per plan tier, per-key override — and distribute rule changes to the fleet without a deploy (hot reload), with a shadow/dark-launch mode.
  • Degrade safely when the counter store is unreachable: fail open onto a per-instance local bucket rather than reject all traffic.
  • (Explicitly out of scope, mentioned not built) concurrency limiting, fleet load-shedding by criticality, and billing/quota accounting — different mechanisms for different problems.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Latency: the added rate-limit check is < 1 ms at p50 with a bounded p99; it must not dominate the latency of the work it protects. A local fast path keeps the tail off the store's tail.
  • Availability: the limiter must be more available than the API it fronts. It fails open — a counter-store outage degrades enforcement, it does not take down the API.
  • Accuracy: enforce ≈ R with bounded, one-sided error. Over-admission during failover / sync lag is acceptable and quantified; false rejection of in-quota traffic is not.
  • Correctness under concurrency: the check is a read-modify-write and MUST be atomic — no interleaving lets two requests both consume the last token.
  • Scalability: horizontal on every tier. Aggregate throughput scales with shards; the residual limit is a single hot key (one slot, one core).
  • Consistency: the per-key counter is linearizable on its primary (atomic Lua, read+write the primary). Rule propagation is eventually consistent (advisory ceilings, not transactions).

04Capacity estimation

How much load and data does this have to hold?

Inputs: 1M req/s peak, 300k avg, 5M active keys, 2% denied, 200 enforcers, 100k ops/s per Redis shard, 64 B/key, 0.5 ms check RTT.

  • Checks/sec: one atomic check per request → 1,000,000 ops/s peak against the counter store. This, not storage, is the sizing driver.
  • Shards for throughput: 1,000,000 / 100,000 = 10 shards minimum; provision 16 for headroom and failover margin. (A single shard sustains ~100k–200k un-pipelined atomic ops/s; the check is synchronous so we do NOT get pipelining's ~1M/s.)
  • Counter memory: 5,000,000 keys × 64 B = 0.32 GB across the cluster — trivial. Token bucket's O(1) state is why: two numbers per key, versus a sliding-window log's one-timestamp-per-request (Figma measured ~20 MB for 10k users × 500 req before switching away from the log).
  • Denied / allowed split: at 2% deny, ~20,000 req/s get a 429 at the gateway and ~980,000 req/s are forwarded — the 20k never reach the upstream, which is the point.
  • Fan-out over-admission: with 200 enforcers and NO coordination (each enforcing the full R locally), the effective ceiling is 200 × R. This single number is why a shared counter exists.
  • Hot-key ceiling: a single key is one hash slot on one shard on one core → capped at ~100,000 ops/s no matter how many shards you add. Aggregate cluster throughput is unreachable for one key.
  • Latency budget: the same-DC check RTT (~0.5 ms) lands on every request's p50 and p99. A cross-region check would be ~50–150 ms — categorically off the synchronous path (hence per-region enforcement + async reconciliation for true-global limits).

05API design

What does the outside world call, and what comes back?

**The check is server-internal; what the client sees is a status + headers.** The gateway calls the limiter over gRPC (check(domain, descriptor)), gets back {allowed, limit, remaining, retry_after_s, reset_s}, and translates that into the HTTP contract below.

Allowed request:

GET /v1/search?q=... HTTP/2
Authorization: Bearer k_live_abc123

200 OK
RateLimit: "default";r=87;t=42          # 87 requests remaining; quota resets in 42s
RateLimit-Policy: "default";q=100;w=60  # policy: 100 requests per 60s window

Denied request (over limit):

GET /v1/search?q=... HTTP/2
Authorization: Bearer k_live_abc123

429 Too Many Requests
Retry-After: 8                          # come back in 8 seconds (takes precedence)
RateLimit: "default";r=0;t=8
RateLimit-Policy: "default";q=100;w=60
Content-Type: application/json

{ "error": "rate_limited", "message": "API rate limit exceeded for this key." }
  • 429 is RFC 6585. It MUST NOT be stored by a cache (so a shared cache never replays a throttle verdict), and it MAY carry Retry-After (RFC 9110) as either delay-seconds (8) or an HTTP-date.
  • RateLimit / RateLimit-Policy are the current IETF fields (draft-ietf-httpapi-ratelimit-headers): RateLimit carries r (remaining) and t (seconds to reset); RateLimit-Policy carries q (quota) and w (window seconds). They collapsed the older RateLimit-Limit/-Remaining/-Reset triplet. When both RateLimit and Retry-After are present, Retry-After takes precedence. The server MAY send RateLimit on success responses too, so clients self-pace before hitting the wall.
  • The de-facto **X-RateLimit-* convention is intentionally avoided: its Reset is ambiguous (GitHub emits an absolute epoch; others emit a delta), which the IETF t (a delta) fixes.

Rule / descriptor model (Envoy-style, what the limiter matches on):

domain: public-api
descriptors:
  - key: tier            # most-specific-match wins
    value: premium
    rate_limit: { unit: minute, requests_per_unit: 1000, algorithm: token_bucket, burst: 200 }
  - key: tier
    value: free
    descriptors:
      - key: route
        value: /v1/search
        rate_limit: { unit: minute, requests_per_unit: 30, algorithm: token_bucket, burst: 10 }
  - key: api_key         # per-key override (emergency clamp / whitelist)
    value: k_live_abc123
    rate_limit: { unit: minute, requests_per_unit: 5, shadow_mode: true }

06Data model

What gets stored, and what is it looked up by?

The counter store holds one value per key — the whole design is O(1) state per key, which is why 5M keys fit in ~0.3 GB.

Token bucket (canonical) — the value packed per key:

fieldtypenotes
tokensfloatcurrent tokens; refilled lazily on read
last_refill_tsint (ms)last time tokens were recomputed
(implicit) keystringrl:{api_key:k_abc}:/v1/search — {...} hash-tag colocates

The Lua script (executed atomically inside Redis) — read, refill, decide, write, all with no interleaving:

-- KEYS[1] = bucket key   ARGV = capacity, refill_rate, now_ms, cost
local b   = redis.call('HMGET', KEYS[1], 'tokens', 'ts')
local tok = tonumber(b[1]) or ARGV[1]            -- new key seeds full
local ts  = tonumber(b[2]) or ARGV[3]
local filled = math.min(ARGV[1], tok + (ARGV[3]-ts)/1000 * ARGV[2])  -- lazy refill
local ok  = filled >= ARGV[4]
if ok then filled = filled - ARGV[4] end
redis.call('HSET', KEYS[1], 'tokens', filled, 'ts', ARGV[3])
redis.call('PEXPIRE', KEYS[1], 120000)           -- idle-key cleanup only, not correctness
return { ok and 1 or 0, filled }                 -- allowed, remaining

Why this shape:

  • Two numbers per key (tokens, ts) → O(1) memory; refill is computed from the timestamp delta, so there's no background dripper.
  • The whole read-modify-write is one script → atomic. The naive GET tokens; if ok: SET splits it across two round trips and races (two requests both read the last token).
  • Read AND write the primary — never a replica. Replicas expire keys lazily (the primary drives expiry via a synthetic DEL), so a replica read can report a window the primary has already rolled over (GitHub's contradictory-header bug: a rejection stapled to r=5000). For long windows, store an explicit expires_at in the value and set the Redis TTL a second later — TTL is cleanup, not correctness.

Rules live in the Control Plane (versioned config), pushed to the limiter's in-memory cache — a descriptor tree matched most-specific-first. Not in the counter store; policy and counters are separate failure domains.

07High-level design

Which components handle a request, and in what order?

The request path (allow): Client → Edge LB (consistent-hash on key) → API Gateway. The gateway builds the descriptor and calls the Rate Limiter Service (gRPC). The limiter resolves the rule and runs the token-bucket Lua against the Redis primary for that key's shard → allowed. The gateway forwards the request to the Protected Upstream and stamps RateLimit headers on the response.

The request path (deny): Same up to the limiter → Redis EVAL returns "no tokens" → the gateway returns 429 immediately with Retry-After; the upstream is never called. The value of edge enforcement is exactly this: the expensive work never starts for over-budget traffic.

Two tiers (the latency + availability insurance): The gateway runs a cheap in-process token bucket in front of the global check. It sheds obvious floods at line rate (protecting the shared limiter), and — critically — when the limiter/store is unreachable it fails open onto that local bucket rather than 429ing everyone. Consistent-hash routing at the LB is what makes the degraded local decision ≈ R instead of 200 × R (a key's whole stream lands on one instance).

The control path: Control Plane → (push rules) → Rate Limiter, and → (push local-limit + shed config) → Gateways. Hot reload / xDS, eventually consistent, shadow_mode for dark launch. Rules are advisory ceilings, so a non-atomic cutover is fine.

The distributed spine in one line: 200 enforcers converge on one global count per key because the counter is a single authoritative shard-owned value updated by an atomic Lua script; everything else (local fallback, consistent-hash routing, fail-open) is about keeping that convergence cheap, fast, and survivable.

08Deep dives

Which part breaks first, and what do you do about it?

1. Which algorithm, and the fixed-window trap. Token bucket is the default: state is (tokens, last_refill_ts), refill is tokens = min(capacity, tokens + elapsed × refill_rate), a request is allowed iff tokens ≥ cost. O(1) memory, one atomic op, and it models real bursty clients (a page load fires 30 calls; the bucket absorbs the burst, then throttles to the sustained rate). The instructive counter-example is the fixed-window counter: INCR a per-window key with a TTL. Cheap, but it admits up to 2×L across a boundary — L requests in the last instant of window N plus L in the first instant of N+1, i.e. 2L in a near-instant straddling the reset (an exact upper bound, not a rough double). The sliding-window log fixes accuracy exactly but costs O(L) memory (a timestamp per request) — too heavy at fleet scale. The sliding-window counter is the pragmatic middle: keep the previous and current window counts and blend them — estimated = prev × ((W − elapsed)/W) + curr. Cloudflare measured this over 400M requests from 270k sources: 0.003% of requests wrongly decided, 6% average rate error, 0 false positives against legitimate traffic. Figma reached for the same thing (60 sub-windows, ~88% less memory than the log). Rule of thumb: token bucket unless you specifically need "no rolling W-window exceeds L" — then sliding-window counter.

2. Atomicity — why Lua, not GET-then-INCR. A rate-limit decision is a read-modify-write. Split across two round trips it races: req A GET → 9, req B GET → 9 (limit 10), both decide "ok", both SET 10 — one request leaked past the limit, and this is worst exactly under the concurrency the limiter exists to control. The gap is in application code, not Redis — a single Redis instance doesn't save you. Redis runs a whole Lua/EVAL script atomically (single-threaded command execution, no interleaving), so read-decide-decrement-set-TTL is one indivisible step. Caveat worth stating: "atomic" here means isolation + serialized execution, not transactional rollback (like MULTI/EXEC, Redis won't undo a write if a later line errors) — fine for a single-pass token-bucket script. INCR alone is atomic for a pure counter, but the moment you need "check limit, maybe set TTL, compute retry-after" you're back to multi-step and need the script.

3. Local vs global — the fan-out that multiplies your limit. This is the heart of "distributed". If each of N enforcers holds its own bucket and enforces the full R against only the traffic it sees, a client spread evenly across them achieves up to N × R (Envoy's own docs: 3 replicas × 50/s = 150/s). Three mitigations, each a trade: per-replica budget R/N is coordination-free but assumes uniform load — under skew a hot instance rejects at R/N while the fleet is far under R (false rejects); shrink the safety factor and you under-utilize. Sticky / consistent-hash routing pins a key to one enforcer so its local count is the global count (Cloudflare's anycast: one IP → one PoP) — but routing churn resets a key's counter and hot keys become hot instances. Token leasing grabs R/N tokens from the central store and burns them locally, round-tripping only to refill — bounds coordination while keeping decisions local. Our canonical picks the fourth option — one authoritative shared counter — and uses consistent-hash routing only to make the fallback correct.

4. The failover free-burst. Redis replication is async: a write is acked to the client, then shipped to replicas. If the primary dies in that window and a replica is promoted, the un-shipped increments are gone — the counter regresses, and the client's budget effectively resets right in the middle of a ~12 s failover, the worst possible moment for the downstream. AOF everysec bounds crash loss to ≤1 s of increments (RDB or "none" would lose minutes / everything — a much bigger burst). WAIT n timeout blocks the check until N replicas ack — it narrows the window but doesn't close it and it taxes every request with latency, so it's rarely worth it. The honest posture: short-window limits tolerate a ≤1 s reset; long-window quotas need the explicit-expires_at-in-value discipline so a promoted replica can still tell the window closed.

5. Hot key = hot shard = hot core. Sharding spreads keys evenly; it does nothing for traffic skew. A single key (a celebrity tenant, an attacker on one API key, or a badly-designed global {} key) is one hash slot (CRC16(key) mod 16384) on one shard on one core — Redis is single-threaded per node, so all of that key's traffic lands on one core while the other 15 shards idle. Adding nodes does not help — you can't split a key across slots, and rebalancing just relocates the hotspot. Fixes: shard the counter into N sub-keys (ctr:{k}:0..N-1), increment a random one, sum on read (loses exactness/atomicity of a single key); local pre-aggregation (each enforcer counts in memory, flushes a batched delta) to cut per-request hits on the hot key; or a coarse local limiter that absorbs the obvious abuse before it ever reaches the shared key. Hash tags help colocation but worsen this if overused (a {tenant} tag piles a whole tenant onto one node).

6. Fail-open vs fail-closed — don't become the outage. When the counter store is unreachable you must choose. Fail-open (allow): brief loss of enforcement / possible abuse burst, but the API stays up — Stripe's explicit stance ("if Redis were to go down, requests wouldn't be affected"; exceptions caught at every level). Fail-closed (reject): the limit is never exceeded, but a Redis blip 429s 100% of traffic — the limiter becomes a self-inflicted, fleet-wide DoS, the classic amplification where a dependency meant to protect availability destroys it. Envoy makes it a flag (failure_mode_deny). Default fail-open, choose fail-closed only when exceeding R is catastrophic (a fragile downstream, hard billing/fraud caps), and always back it with a local fallback bucket so fail-open doesn't mean no limit. The deeper answer (Google SRE) is to not centralize correctness at all: layer client-side adaptive throttling + per-service self-protection so a degraded limiter degrades gracefully.

7. One global limit across regions. A single global counter can't live near every enforcer. A synchronous cross-region read is ~50–150 ms — off the table for a per-request check. So a true global limit means per-region local decisions + async cross-region reconciliation (eventual consistency; a guaranteed over-admission window ≈ replication lag), or a strongly-consistent multi-region store (DynamoDB global tables MRSC) at real $/latency cost — replicated-write and storage cost scale ~linearly with region count. The pragmatic answer most teams pick: don't enforce one global limit — split R into per-region sub-limits (the per-replica-budget idea one level up), accept bounded global over-admission, and reserve a true global counter for the rare case where exceeding R is worth cross-region coordination.

8. The HTTP contract and the retry storm. Return 429 (RFC 6585, never cached) with Retry-After and the IETF RateLimit/RateLimit-Policy fields so clients know their standing and their reset. But the server half is only half the contract: if clients retry immediately and in lockstep, a 429 becomes the thundering herd it was meant to prevent (Cloudflare's Nov-2023 DR recovery: synchronized reconnects overwhelmed services until limits throttled the surge). Clients MUST honor Retry-After, back off with jitter, and ideally run adaptive throttling (p(reject)=max(0,(requests−2·accepts)/(requests+1))) to shed locally before hitting the network — because even rejecting a request costs the server CPU (RFC 6585 says so, and it's why you may drop the connection instead of 429ing under a flood).

09Trade-offs

What did this design cost, and what breaks at 10×?

What we accepted:

  • Approximate global enforcement (AP over CP). The limiter is fail-open and tolerates bounded over-admission (failover free-burst, sync lag, fixed/sliding-window slack). The alternative — block every request on a strongly-consistent quorum — makes the store's availability the API's availability. An over-admitted request is cheap; a wrongly-rejected legitimate request is expensive. This is the load-bearing product decision.
  • Central-authoritative counter + local fallback, not local-first async reconciliation. We pay ~0.5 ms per request for an accurate global decision (the local tier keeps the tail off the store), and degrade to the local bucket only on store failure. The Cloudflare-style local-first / sloppy-counter model would be lower-latency and more available but always approximate — we chose accuracy in steady state, approximation only when degraded.
  • Token bucket over sliding-window counter as the default. Burst tolerance + O(1) + one atomic op. We give up the "no rolling W-window exceeds L" guarantee the sliding-window counter provides; for the common case (protect a downstream, be fair between tenants) burst-friendly is the better fit, and we document sliding-window counter as the swap-in when exactness matters.
  • AOF everysec (≤1 s counter loss on crash), read+write the primary. We accept a ≤1 s free-burst on a crash and 2× read cost (no replica offload) in exchange for correct, non-stale counters. AOF always (per-write fsync) would close the window but ~10× the write cost on the hottest path; replica reads would be faster but wrong (lazy expiry).
  • Single region. A regional outage takes the limiter with the API it fronts — acceptable because the limiter fails open, so "limiter region down" degrades to "no enforcement", not "no API". True multi-region global limits are a documented deep-dive, not the canonical.
  • Consistent-hash routing's costs — hot keys become hot instances, and a fleet reshuffle briefly resets a key's local counter — accepted to make the fail-open fallback ≈ R instead of 200 × R.

What breaks at 10× scale (~10M req/s):

  • Aggregate check throughput — go from 16 to ~100+ Redis shards; the check-per-request model is the cost driver, so consider a local-aggregation tier (flush batched deltas) to cut per-request hits.
  • A single hot key still can't exceed one shard — at 10× the abusive-key ceiling bites sooner; sharded counters / local pre-aggregation become mandatory, not optional.
  • The gateway→limiter gRPC hop — collapse the limiter into an Envoy filter (in-process) to drop a network hop from the tail, or move to redis-cell/GCRA to shave Lua overhead.
  • Multi-region becomes unavoidable — accept per-region sub-limits with async reconciliation; a true global limit stays reserved for the few keys where it's worth the cross-region cost.

Per-component failure stories: see "Failure scenarios we model" below.

Trace catalogue

The simulator models 4 rate-limiter journeys, each a deterministic walk through the canonical hldNodes/hldEdges with budgets from the SLO table above:

Trace IDJourneyBudget
rate-limiter:allow-under-limitUnder limit: global token-bucket check passes → request forwarded to the upstream60ms
rate-limiter:deny-over-limitOver limit: check fails → 429 at the gateway, upstream never called10ms
rate-limiter:store-fail-openCounter store unreachable → gateway fails OPEN on its local bucket (degraded)30ms
rate-limiter:rule-propagationControl plane pushes a limit change to the limiter and the enforcers (hot reload)5000ms

Each is satisfied against the user's actual graph; gaps surface as concrete prompts ("you have no store behind your rate limiter — your check can't be global, so every enforcer limits locally and an attacker gets N× the limit").

Failure scenarios we model

The simulator ships 5 chaos scenarios specific to a distributed rate limiter, each grounded in a cited real-world failure mode:

Scenario IDReal-world precedentCategorySeverity
rate-limiter:counter-store-failoverRedis async-replication failover (primary loss) → un-replicated increments lost → free burst; GitHub sharded-Redisdatahigh
rate-limiter:hot-key-shardRedis hot-key on one slot/core (Percona/Redis hot-shard docs; Figma hot-key pruning)traffichigh
rate-limiter:store-unreachable-failopenCounter store partition → fail-open vs fail-closed (Stripe fail-open; Envoy failure_mode_deny)infrahigh
rate-limiter:limiter-tier-overloadDecision tier saturates → check latency on every request; two-tier local-fast-path (Envoy local+global)trafficmedium
rate-limiter:control-plane-downConfig/control-plane outage or bad-rule push (Cloudflare 2019 regex-CPU; Knight Capital 2012 partial deploy)processmedium

Each scenario mutates a specific edge or node according to the canonical's authored internals (e.g. counterstore.failover.rtoSeconds=12 / rpoSeconds=1 → the failover applies a 12-second decision-degraded window and a ≤1-second counter-loss free-burst) and returns a diff explaining why a probe degraded or broke.

Primary sources

  • Stripe — Scaling your API with rate limiters (token bucket on Redis, fail-open)
  • Cloudflare — How we built rate limiting to millions of domains (sliding-window counter)
  • GitHub — Sharded, replicated rate limiter in Redis (replica-expiry gotcha)
  • Envoy — Global rate limiting + Lyft `ratelimit` service (local + global)
  • Figma — An alternative approach to rate limiting (sliding-window counter, hot key)
  • Google SRE Book — Handling Overload (adaptive throttling, criticality)
  • IETF draft-ietf-httpapi-ratelimit-headers (RateLimit / RateLimit-Policy)
  • redis-cell — GCRA rate limiting as one command (CL.THROTTLE)

Now defend it

Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.

Work Distributed Rate Limiter yourself