Rate limiter: one limit for a fleet of gateways
Cap how many requests each client makes, in front of an API that runs on many servers. The hard part is one limit across many gateways, and what to do when the counter store fails.
- 1Decide in the gateway, count in Redis: one atomic script per request gives every gateway one exact limit.
- 2Memory-only limits break when the fleet scales: N gateways admit up to N times the limit.
- 3Local counts with a periodic sync cut Redis load, and over-admit by about one sync interval of traffic.
- 4Rules live in Postgres with a version; every gateway reloads them without a restart.
"Design a rate limiter that caps how many requests a client may make, in front of an API served by many servers."
| question | answer assumed |
|---|---|
| Who is a client? | An API key; an IP for anonymous calls |
| What traffic? | 5 billion calls a day, 10 million keys |
| Limits by what? | Tier and route: a few rules per request |
| How exact? | Within a few percent; never 2× over |
| One region or many? | One first; many as a follow-up |
| Billing quotas too? | No; quotas are durable and separate |
- functional
- Allow or refuse each call. Tell the client when to retry. Change rules without a deploy.
- latency
- Under 1 ms added at p50: one Redis round trip at most.
- availability
- A limiter failure must not take the API down.
- accuracy
- Exact for login and paid calls; approximate is fine for reads.
I assume 5 billion API calls a day, 10 million clients, 2 rules per request, and under 1 ms added at p50.
| number | arithmetic | result |
|---|---|---|
| calls, average | 5 B ÷ 86,400 s | 57,870/s |
| calls, peak (3×) | 57,870 × 3 | 173,611/s |
| decisions, peak | 173,611 × 2 rules | 347,222/s |
| Redis nodes, shared count | 347,222 ÷ 25,000 | 14 shards |
| memory per client and rule | sliding counter, measured | 142 B |
| memory, all clients | 10 M × 2 × 142 B | 2.8 GB |
| per gateway, 50 gateways | 173,611 ÷ 50 | 3,472/s |
5 billion calls, 10 million keys, 2 rules, a 3× peak and 50 gateways are assumptions. 25,000 script decisions a second per Redis is the B3 lab number. Memory per key is Redis MEMORY USAGE, measured here on Redis 8.10.2.
At peak the limiter makes about 347,222 decisions a second. Memory is small; throughput needs a Redis Cluster or fewer Redis calls.
| call | body | success | errors | notes |
|---|---|---|---|---|
| check(client, route, cost) | in-process, or one RPC | {allowed, remaining, retry_after_ms} | store down: the rule's fail mode | Called by the gateway on every request |
| any API call, allowed | unchanged | 2xx + RateLimit-Remaining | The client can slow down before it is refused | |
| any API call, refused | none | 429 + Retry-After: 8 | Whole seconds, rounded up | |
| any API call, limiter down | none | 503 + Retry-After: 1 | Fail-closed rules only | |
| PUT /admin/rules/{tier}/{route} | {limit, window_s, fail_open} | 200 {version} | 400 bad rule, 409 no default | PUT is idempotent |
- 429 is defined in RFC 6585 and Retry-After in RFC 9110. An IETF draft adds RateLimit and RateLimit-Policy header fields; until it is final, many APIs send X-RateLimit-* headers.
- Key the limit by the authenticated client, not the IP. Many users share one NAT address.
- A request with a cost, such as a batch of 50, takes 50 from the count.
The limiter has no public endpoint of its own. Clients see 429 with Retry-After; operators change rules through an admin call that bumps a version.
- rules
- Tens of rows; read once per change, never per request
- counters
- One key per client, rule and window; about 20 million at peak
- hash tag
- {client}: every key of one client sits in one Redis Cluster slot
-- One row per limit. '*' matches any tier or any route.
CREATE TABLE rate_rules (
tier text NOT NULL,
route text NOT NULL,
limit_count int NOT NULL CHECK (limit_count > 0)2,
window_ms int NOT NULL CHECK (window_ms >= 1000),
fail_open boolean3 NOT NULL DEFAULT true,
PRIMARY KEY (tier, route)1
);
-- One version for the whole rule set. Every change to rate_rules bumps it.
CREATE TABLE rule_version (
id int PRIMARY KEY CHECK (id = 1),
version bigint NOT NULL
);
INSERT INTO rule_version VALUES (1, 1);
CREATE FUNCTION bump_rule_version() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
UPDATE rule_version SET version = version + 1 WHERE id = 1;
RETURN NULL;
END $$;
CREATE TRIGGER rate_rules_changed
AFTER INSERT OR UPDATE OR DELETE ON rate_rules
FOR EACH STATEMENT4 EXECUTE FUNCTION bump_rule_version();- 1One rule per tier and route. '*' matches any; the most specific rule wins.
- 2The database refuses a rule that would block everything by mistake.
- 3The failure mode is part of the rule: open for public reads, closed for login.
- 4One bump per change, in the same transaction. A gateway never sees new rules with an old version.
Rules are a small versioned table in Postgres. Counters are small expiring keys in Redis, one per client and window, under a hash tag.
Step 1: Check
- Find the client from its API key. Read its rule for this route from memory.
- Run the limiter script on the client's key: read, decide and write as one step.
- Allowed: forward the request, with the remaining count in a header.
If it fails
Redis is slow: a short timeout, then the rule's fail mode. The API never waits on the limiter.
The gateway decides; Redis counts; Postgres holds the rules. The API servers see only the calls the limiter admitted.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Redis | Server-side scripts (EVALSHA) | Read the count, decide and write as one step. Every gateway shares one exact limit. | Inventory holds, idempotency |
| Redis | INCR, INCRBY and PEXPIRE | A counter per window that deletes itself. INCRBY pushes a batch of local admits in one command. | Page views, quotas |
| Redis | Pipelining | A sync sends one INCRBY per key in one round trip. | Bulk loads |
| Redis | Cluster hash tags: rl:{client} | All keys of a client in one slot, so one script can use several windows. | Any multi-key script |
| Redis | Asynchronous replication | Limit A failover can lose recent counts. A client may get a little extra for one window. | |
| Postgres | Constraints and a statement trigger | Bad rules are refused; every change bumps one version, in the same transaction. | Config tables, feature flags |
| Postgres | One-statement snapshot | The version and the rules come from one snapshot, so they always match. | Consistent reports |
| Service | In-memory counters and an atomic pointer | The local tier and the sync design; rules swap with no lock on the request path. | Feature flags, config |
| HTTP | 429, Retry-After, RateLimit-* headers | A standard refusal that tells the client how long to wait. | 503 during maintenance |
Redis gives me atomic scripts, counters with expiry and hash tags. Postgres gives me a durable, versioned rule set. The gateway gives me one place to enforce both.
| where | runs on | status |
|---|---|---|
| Gateway, in process, counts in Redis | Redis | Approved |
| Sidecar per service, counts in Redis | Redis | Service mesh |
| Library in each service, counts in Redis | Service | One language |
| A separate limiter service, by RPC | Service | One more hop |
| Edge or WAF, by IP | Edge | IP floods only |
| Each service, in memory only | Service | N × the limit |
on request(client, route):
rule = rules.for(client.tier, route) // memory, no I/O1
d = redis_script(key(client, rule), now) // ONE atomic step2
IF d.allowed:
forward; add RateLimit-Remaining3: d.remaining
ELSE:
RETURN 429, Retry-After: ceil(d.retry_after)4- 1Rules are loaded into memory and reloaded on a version change, never read per request.
- 2Redis runs the script to the end before any other command. 200 gateways at once still admit exactly the limit.
- 3A well-behaved client slows down before it is refused.
- 4Retry-After takes whole seconds. 7.7 s becomes 8; rounding down invites a second refusal.
Tested source Redis script: fixed window · Go: refuse with 429
-- Fixed window: count requests in the current window; refuse once the count passes the limit.
-- KEYS[1] the counter for this window: rl:{user:42}:fw:<window start, ms>
-- ARGV[1] now, ms ARGV[2] window, ms ARGV[3] limit
-- Returns {allowed, retry after ms, remaining, delay ms}.
local now = tonumber(ARGV[1])
local window = tonumber(ARGV[2])
local limit = tonumber(ARGV[3])
local count = redis.call('INCR', KEYS[1])
if count == 1 then
redis.call('PEXPIRE', KEYS[1], window)
end
if count <= limit then
return {1, 0, limit - count, 0}
end
local window_end = now - now % window + window
return {0, window_end - now, 0, 0}if !dec.Allowed {
secs := (dec.RetryAfter + time.Second - 1) / time.Second
w.Header().Set("Retry-After", strconv.FormatInt(int64(secs), 10))
http.Error(w, "too many requests", http.StatusTooManyRequests)
return
}| algorithm | Redis state | status |
|---|---|---|
| Fixed window | 68 B | 2× at a boundary |
| Sliding window counter | 142 B | Close to exact |
| Token bucket | 86 B | Bursts |
| Sliding log | 1 entry per call | Small limits |
The five algorithms, their scripts and a recorded burst are on B3. This sheet uses the fixed window, so the fleet experiment counts one key per window.
I put the limiter in the gateway, in process, with the counts in Redis. In a service mesh, a sidecar does the same job per service.
| design | where | status |
|---|---|---|
| Each gateway allows the full limit, in memory | Service | N × the limit |
| Each gateway allows limit ÷ N, in memory | Service | Even spread, fixed N |
| One script per request on a shared count | Redis | Exact |
| Local counts, synced every interval | Redis | Approximate |
| Lease tokens from Redis in batches of 10 | Redis | Approximate |
| Sticky routing: one client, one gateway | Balancer | Until it moves |
| 16 gateways, lab | busiest window | fair client lost | Redis cmds/s |
|---|---|---|---|
| Local, full limit | 4,000 | 0 | 0 |
| Local, limit ÷ 16 | 992 | 38 | 0 |
| Shared Redis | 1,000 | 0 | 490 |
| Local + sync, 1 s | 1,395 | 0 | 55 |
| Local + sync, 0.25 s | 1,098 | 0 | 194 |
Limit 1,000 per 10 s. A heavy client sends 400 a second; a fair client 90. Sync at 1 s uses 9× fewer Redis commands than the shared count and lets 40% extra through in the busiest window.
allow(client, now): // memory only1
w = window of now
IF known[client, w] + pending[client, w]2 >= limit:
RETURN refuse
pending[client, w] += 1
RETURN allow
every interval: // one round trip3
FOR EACH (client, w) seen in this window:
known[client, w] = INCRBY global(client, w), pending4
pending[client, w] = 0
IF Redis fails: keep pending for the next sync5- 1No network call on the request path. Redis load no longer grows with traffic.
- 2The global count at the last sync, plus what this gateway admitted since. Other gateways’ recent admits are missing: that is the over-admission.
- 3A pipeline: one INCRBY per key this gateway saw. Cost grows with gateways and keys, not with requests.
- 4Adds this gateway’s admits and returns the total from every gateway in one command.
- 5A failed sync loses no counts. The next one pushes them.
Tested source Go: allow and sync
// Allow implements Gate. It never calls Redis.
func (g *SyncGate) Allow(_ context.Context, client string, now time.Time) (bool, error) {
k := winKey{client, windowStart(now, g.Window)}
g.mu.Lock()
defer g.mu.Unlock()
if g.known[k]+g.pending[k] >= g.Limit {
return false, nil
}
g.pending[k]++
return true, nil
}
// Sync adds this instance's pending admits to the global counts and reads them back. It returns
// the number of Redis commands it sent.
func (g *SyncGate) Sync(ctx context.Context, now time.Time) (int, error) {
cur := windowStart(now, g.Window)
g.mu.Lock()
batch := map[winKey]int64{}
for k := range g.known {
if k.start < cur {
delete(g.known, k) // a past window can no longer admit anything
} else {
batch[k] = 0 // read the global count even if this instance admitted nothing
}
}
for k, n := range g.pending {
batch[k] = n
}
g.pending = map[winKey]int64{}
g.mu.Unlock()
if len(batch) == 0 {
return 0, nil
}
keys := make([]winKey, 0, len(batch))
cmds := make([]*redis.IntCmd, 0, len(batch))
_, err := g.RDB.Pipelined(ctx, func(p redis.Pipeliner) error {
for k, n := range batch {
keys = append(keys, k)
cmds = append(cmds, p.IncrBy(ctx, SyncKey(k.client, k.start), n))
if n > 0 {
p.PExpire(ctx, SyncKey(k.client, k.start), 2*g.Window)
}
}
return nil
})
if err != nil {
g.mu.Lock()
for k, n := range batch {
g.pending[k] += n // keep the admits; the next sync pushes them
}
g.mu.Unlock()
return 0, fmt.Errorf("sync %d counters: %w", len(batch), err)
}
g.mu.Lock()
sent := 0
for i, k := range keys {
g.known[k] = cmds[i].Val()
sent++
if batch[k] > 0 {
sent++ // the PEXPIRE
}
}
g.mu.Unlock()
return sent, nil
}
A shared Redis count is exact and costs one call per request. For the heaviest keys, I count locally and sync every second; that over-admits by about one second of traffic.
Each gateway decides from memory and pushes its counts to Redis at an interval. Limit: 1,000 requests per 10 s per client. A heavy client sends 400 a second; a fair client sends 90. A random balancer spreads both over 8 gateways.
- heavy client got, 30 s
- 4,050 exact limit: 3,000 (+1,050)
- busiest 10 s window
- 1,359 limit 1,000
- fair client refused
- 0 of 2,700, all under the limit
- Redis commands a second
- 28 8 round trips a second
- admitted
- refused, 429
| 8 gateways | heavy got (exact 3,000) | busiest window | fair refused | Redis commands/s |
|---|---|---|---|---|
| Local, full limit | 12,000 | 4,000 | 0 | 0 |
| Local, limit ÷ N | 3,000 | 1,000 | 3 | 0 |
| Shared Redis | 3,000 | 1,000 | 0 | 490 |
| Local + sync, 1 s | 4,050 | 1,359 | 0 | 28 |
With memory-only limits, more gateways means more admits. The shared count stays exact at any fleet size; the synced design stays close, with far fewer Redis calls.
| mode | clients see | status |
|---|---|---|
| Fail open, local cap on | 200 up to the local cap | Public reads |
| Fail closed | 503 + Retry-After | Login, paid calls |
| Fail open, no cap | 200, no limit at all | Short outages |
| Wait for Redis, no timeout | Every request hangs | Not approved |
d = limiter(client, now) // with a short timeout1
IF the limiter fails:
IF rule.fail_open2: log it; serve the request
ELSE: RETURN 5033, Retry-After: 1
IF NOT d.allowed:
RETURN 429, Retry-After: whole seconds, rounded up- 1Well below the route’s latency budget. A slow limiter must not slow every request.
- 2Stored with the rule, so the choice is made in advance and reviewed like any other change.
- 3503, not 429: the client did nothing wrong.
Tested source Go: fail open or closed · Go: refuse with 429
if err != nil {
if mode == FailOpen {
log.Warn("rate limiter unavailable; failing open", "client", c, "err", err)
next.ServeHTTP(w, r)
return
}
log.Error("rate limiter unavailable; failing closed", "client", c, "err", err)
w.Header().Set("Retry-After", "1")
http.Error(w, "rate limiter unavailable", http.StatusServiceUnavailable)
return
}if !dec.Allowed {
secs := (dec.RetryAfter + time.Second - 1) / time.Second
w.Header().Set("Retry-After", strconv.FormatInt(int64(secs), 10))
http.Error(w, "too many requests", http.StatusTooManyRequests)
return
}| hot key | fix |
|---|---|
| One tenant sends a large share of calls | Local tier or sync for that key |
| One key per tenant for every route | Split the key per route |
| A flood from one client | Local cap refuses it in memory |
Each rule carries its failure mode. A hot tenant gets a local tier, so it stops costing one Redis call per request.
rule_for(tier, route): // most specific wins1
try (tier, route), (tier, *), (*, route), (*, *)
every 10 s, on each gateway:
v = SELECT version // one tiny read2
IF v == current.version: done
set = SELECT every rule // one snapshot3
IF set has no (*, *) default:
keep the old set; alert4
current = set // one atomic swap5- 1A rule for this tier and route beats a tier rule, which beats a route rule, which beats the default.
- 2Polling the version costs one row read per gateway per interval, whatever the number of rules.
- 3The version and the rules come from one statement, so they always match.
- 4A bad change never removes every limit. Tested: a set without a default is refused.
- 5Requests read either the old set or the new one, never a mix. Tested with 8 readers during a reload.
Tested source Go: match a rule · Go: reload
// For returns the most specific rule: tier and route, then tier only, then route only, then the
// default.
func (rs *RuleSet) For(tier, route string) Rule {
for _, k := range [][2]string{{tier, route}, {tier, "*"}, {"*", route}, {"*", "*"}} {
if r, ok := rs.rules[k]; ok {
return r
}
}
panic("rule set without a default rule") // Reload refuses such a set
}
// Reload reads the version, and only when it changed, every rule. It checks the new set before
// it swaps it in with one atomic store; a bad set leaves the old one in use.
func (r *Rules) Reload(ctx context.Context) (bool, error) {
var v int64
if err := r.db.QueryRow(ctx, sqls["version"]).Scan(&v); err != nil {
return false, fmt.Errorf("read rule version: %w", err)
}
if old := r.cur.Load(); old != nil && old.Version == v {
return false, nil
}
rows, err := r.db.Query(ctx, sqls["load"])
if err != nil {
return false, fmt.Errorf("load rules: %w", err)
}
defer rows.Close()
set := &RuleSet{rules: map[[2]string]Rule{}}
for rows.Next() {
var rl Rule
var ms int
if err := rows.Scan(&set.Version, &rl.Tier, &rl.Route, &rl.Limit, &ms, &rl.FailOpen); err != nil {
return false, fmt.Errorf("scan rule: %w", err)
}
rl.Window = time.Duration(ms) * time.Millisecond
set.rules[[2]string{rl.Tier, rl.Route}] = rl
}
if err := rows.Err(); err != nil {
return false, fmt.Errorf("read rules: %w", err)
}
if _, ok := set.rules[[2]string{"*", "*"}]; !ok {
return false, fmt.Errorf("version %d: %w", v, ErrNoDefault)
}
r.cur.Store(set)
return true, nil
}
| several regions | status |
|---|---|
| One Redis in one region for all | Ocean round trip |
| Each region gets limit ÷ regions | Even traffic |
| Count per region, share counts every few seconds | Approximate |
| A counter CRDT: one count per region, summed | Approximate |
| Route each client to a home region | Exact, if possible |
A round trip across an ocean is about 150 ms (R2). Shared regional counts are the sync design at a larger scale: the error grows with the share interval.
Rules reload on a version change with one atomic swap. Across regions, each region counts locally and shares its counts in the background, so the global limit is approximate.
| event | result | why it stays safe | saved by |
|---|---|---|---|
| Redis goes down | No global count. | Each rule's fail mode applies. Fail open keeps the local cap. | Local tier |
| Redis gets slow | Every check waits. | A short timeout, then the fail mode. The API never waits long on the limiter. | Timeout |
| Redis fails over | Recent counts are lost. | A client gets up to one extra window. Acceptable for a rate limit; not for billing. | Replica |
| The fleet scales from 8 to 16 gateways | Memory-only limits double. | The shared count does not depend on the number of gateways. | Redis |
| A sync fails | Gateways decide on stale totals. | Admits stay in pending and go out with the next sync. Tested with Redis unreachable. | Pending counts |
| An operator deletes the default rule | Some calls would have no limit. | The reload refuses the set and keeps the old rules. | Version check |
| Postgres is down | Rules cannot change. | Gateways keep the rules in memory. Checks never touch Postgres. | Rules in memory |
| Refused clients retry at once | A retry storm. | Retry-After plus client jitter; the local cap refuses the storm in memory. | Retry-After |
| One tenant floods one key | One Redis shard is hot. | A local tier or sync for that key; split the key per route. | Local tier |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | One Redis, one script per request, in the gateway. | About 25,000 exact decisions a second. | Floods of refused calls use Redis capacity. |
| 2 | A local cap in each gateway, above a fair share. | A flood stops in memory. Redis sees only what the cap admits. | Real traffic alone nears one Redis. |
| 3 | Redis Cluster, keys under {client}. | Decisions grow with shards: 14 at the lab rate for 347,222 a second. | A few tenants make most calls, and their shards run hot. |
| 4 | Local counts with a 1 s sync for heavy keys. | A heavy key costs one INCRBY per gateway per second. Lab: 9× fewer commands, 40% over at worst. | Clients in several regions. |
| 5 | Per-region counts, shared in the background. | Every check stays in its region. The global limit is approximate. | Top of the ladder. |
Demand as in panel B. Capacity bars are the B3 lab numbers from one loaded laptop; read them as orders of magnitude.
I start with one Redis and one script per request. I add a cluster when decisions outgrow one node, and local counts for heavy keys when Redis load is the cost.
| question | answer | deeper |
|---|---|---|
| Token bucket or sliding window? | A token bucket allows a burst, then the average rate. A sliding window counter is close to exact for a fixed count per window. | B3 |
| How do you shard the counters? | By client key, with a hash tag so one client stays in one slot. Hashing spreads clients; it cannot split one hot client. | X2 |
| Can a gateway use its own clock? | Clocks drift between machines, so windows disagree. Use the Redis clock, or tolerate a small skew in a window that is many seconds long. | X4 |
| How do regions agree on one count? | They do not, in real time. Each region keeps its own count and merges the others' in the background, like a counter CRDT. | X7, X3 |
| How does the balancer affect local limits? | A random or round-robin balancer spreads one client over every gateway. Sticky routing keeps it on one, until that gateway goes away. | B2 |
| What about a monthly paid quota? | That is billing, not protection. Keep it in a durable database with a ledger, and check it less often. | T9 |
| How do you size the Redis fleet? | Decisions a second at peak ÷ what one node does, then headroom. Memory is rarely the limit. | R2 |
| What does the limiter's own availability need to be? | Less than the API's, because it can fail open. Decide which routes may. | O1 |
0 of 9 known
16 gateways each enforce 1,000 per 10 s in memory. What can one client get?
Then give each gateway 1,000 ÷ 16. What goes wrong?
Why one script in Redis, not GET then INCR from the gateway?
Local counts synced every 1 s. How far over the limit can a client get?
Which Redis key holds a client's counters, and why the braces?
Redis is down. What does the gateway do?
Why 429 and not 503 for a client over its limit?
An operator changes a limit. When does every gateway use it?
One tenant sends 30% of all traffic. What breaks first?
- demand
- 5 B calls a day: 57,870 a second, 173,611 at peak; 347,222 decisions with 2 rules.
- one Redis
- About 25,000 script decisions a second; about 55,000 plain INCRs.
- memory
- 68 B a fixed window, 142 B a sliding counter, 86 B a token bucket.
- local only
- 16 gateways: 4,000 in a window against 1,000.
- sync, 1 s
- 40% over at worst; 9× fewer Redis commands.
- reload
- One version read per gateway every 10 s.
- Retry-After
- Whole seconds, rounded up.
Fleet runs replay 30 s of traffic on a controlled clock against Redis 8.10.2, so they are exact. Throughput is from B3, on one shared laptop.