Limits across many servers.
A limit of 1,000 requests a second means nothing if each of 50 gateway pods counts on its own. Keeping one count means shared state, and shared state brings races, hot keys, outages and cross-region latency. This design keeps the count in a sharded in-memory store, makes each check atomic, answers most checks locally, and decides ahead of time what happens when the counter cannot be reached.
Builds on Rate-limiting algorithms.
Requirements.
150 gateway pods in 3 regions, and one limit per merchant. Which algorithm to run is settled in Algorithms; here the question is where the count lives and how every pod agrees on it. Shiplane and every figure about it are our own assumptions.
Functional requirements
Capacity estimates
Sizing the busiest region
- Peak requests (all regions)
- 300K/sassumption
- us-east share of traffic
- 60%split 60/25/15 across us-east, eu-west, ap-south
- Gateway pods per region
- 50
- Buckets checked per request
- 2the account-wide limit plus the route's limit
- Bucket updates per Redis primary
- 50K/sa conservative planning figure for a short script; benchmark your own
- Memory per bucket
- ~120 Bkey, two fields, hash and TTL overhead; a login log of up to 5 entries is a little larger, which this round figure absorbs
- Live buckets
- 3.24M40,000 accounts × 6 route classes = 240K, plus ~3M login logs: the 2M anonymous IPs an hour that Algorithms assumes plus ~1M usernames tried (assumption; an upper bound, since per-IP logs expire 60 s after the newest attempt)
- Peak requests in us-east300K × 0.6180K/sfrom Peak requests (all regions) and us-east share of traffic
- Bucket updates in us-east180K × 2360K/sfrom Peak requests in us-east and Buckets checked per request · Both buckets share a hash tag, so they go in one script call; the round trips are 180K/s.
- Primaries at 50% headroom360K ÷ (50K × 0.5) = 14.416 primariesfrom Bucket updates in us-east and Bucket updates per Redis primary · Rounded up to 16, so each owns an even 1,024 of the 16,384 hash slots.
- Requests per gateway pod180K ÷ 503.6K/s (7.2K bucket updates/s)from Peak requests in us-east and Gateway pods per region
- Memory for all buckets3.24M × 120 B~390 MBfrom Live buckets and Memory per bucket · Fits in any one node. Memory is not what sizes this store.
- Round trips per second, not memory, size the counter store.
- Every one of those round trips is on a request's critical path, so each must stay inside the region and well under 1 ms.
High-level design.
The data path (gateway to counter store) runs on every request; the control path (rules) does not, so a rules outage never blocks traffic. Each region is a full copy of the picture below.
Distributed rate limiter (one region)
| Case | What happens | Rate the client gets |
|---|---|---|
| Pinned connections | The client keeps 2 long-lived HTTP/2 connections, which the load balancer pins to 2 pods. Only those 2 pods see it. | 2 × 20 = 40/s (4% of 1,000) |
| Spread evenly | Thousands of short connections land on all 50 pods about equally. | ~1,000/s (right by luck) |
| Autoscaling | The fleet grows to 75 pods but each still allows 20/s. | 75 × 20 = 1,500/s (150%) |
| Pod restart | A restarted pod starts with a full bucket, and its clients briefly get an extra burst. | up to +40 at once |
Sharing the count safely
Moving the bucket into a shared store fixes the arithmetic, but a naive read-then-write lets two pods spend the same token.
The lost update
- pod-07 → Counter store: HGET rl:{acct_7Q2}:global tokens
- pod-31 → Counter store: HGET rl:{acct_7Q2}:global tokens
- Counter store → pod-07 (reply): "1"
- Counter store → pod-31 (reply): "1"
- Note over pod-07 and pod-31: both see 1 ≥ 1, both allow
- pod-07 → Counter store: HSET … tokens 0
- pod-31 → Counter store: HSET … tokens 0
- Note over pod-07 and Counter store: two requests admitted on one token
One atomic script
- pod-07 → Counter store: EVALSHA bucket 2 rl:{acct_7Q2}:global rl:{acct_7Q2}:labels …
- Note over Counter store: reads TIME, refills, checks every key, spends, sets PEXPIRE: one step
- Counter store → pod-07 (reply): grant 1, retry 0
- pod-31 → Counter store: EVALSHA bucket 2 … (queued behind pod-07's script)
- Counter store → pod-31 (reply): grant 0, retry 10 ms
Redis runs one script at a time and nothing else runs while it does, so the read, the refill and the spend cannot interleave with another pod's. The same trap exists for fixed-window counters: the Redis docs' INCR pattern warns that a client which dies between INCR and EXPIRE leaves a key that never expires, and fixes it with a script. Both of a request's buckets use the hash tag {acct_7Q2}, so Redis Cluster hashes only that part and puts them in the same slot, which a multi-key script requires. How slots map to primaries is covered in Sharding a cache.
The check, as one script (ours)
-- KEYS: the request's buckets, all tagged {acct} so they share one slot
-- ARGV: want, pods, then rate (tokens/s), burst and cost for each key in order
local want, pods = tonumber(ARGV[1]), tonumber(ARGV[2])
local t = redis.call('TIME') -- the store's clock, never a pod's
local now = t[1] * 1000 + math.floor(t[2] / 1000)
local tok, grant, retry = {}, want, 0
for i, key in ipairs(KEYS) do
local rate = tonumber(ARGV[3 * i]) -- tokens/s = milli-tokens per ms
local burst = tonumber(ARGV[3 * i + 1]) * 1000
local cost = tonumber(ARGV[3 * i + 2]) -- tokens one request spends here
local s = redis.call('HMGET', key, 'tokens_milli', 'ts_ms')
local v = (tonumber(s[1]) or burst) + (now - (tonumber(s[2]) or now)) * rate
tok[i] = math.min(burst, v)
local whole = math.floor(tok[i] / 1000)
if whole < cost then -- too few tokens here: deny all
grant = 0
retry = math.max(retry, math.ceil((cost * 1000 - tok[i]) / rate))
else -- lease: at most half of what is left
grant = math.min(grant, math.max(1, math.floor(whole / (2 * pods * cost))))
end
end
for i, key in ipairs(KEYS) do
local rate = tonumber(ARGV[3 * i])
local burst = tonumber(ARGV[3 * i + 1]) * 1000
local cost = tonumber(ARGV[3 * i + 2])
redis.call('HSET', key, 'tokens_milli', tok[i] - grant * cost * 1000, 'ts_ms', now)
redis.call('PEXPIRE', key, math.ceil(burst / rate) * 2)
end
return { grant, retry } -- grant (requests) 0 means deny
| Line | Why |
|---|---|
| TIME | Pods' clocks drift by milliseconds; one clock per bucket keeps refill arithmetic consistent. Scripts may call TIME because Redis replicates a script's writes, not the script (the default since Redis 5). |
| tokens_milli | Thousandths of a token. For per-second plans it stays an integer (a Free key refills 10 milli-tokens per ms), so there is no rounding drift. A per-minute or per-hour bucket would pass a fractional rate (5/min = 0.083 milli-tokens per ms); Lua numbers are doubles, and the rounding error over a bucket's life is far below one token. |
| all-or-nothing | The first loop decides and the second spends, so a request never uses the account's token when its route bucket is empty. |
| PEXPIRE | A bucket untouched for 2 × B ÷ R is full anyway, so deleting it loses nothing and frees memory. That is 4 s on every plan (Free 2 × 20 ÷ 10, Growth 2 × 200 ÷ 100, Enterprise 2 × 2,000 ÷ 1,000). Login rules are not buckets but sliding logs (see Algorithms), checked by a similar script on a sorted set and expired one window after the newest entry. |
| want and pods | With want 1 it is an exact per-request check. With want 50 it grants a lease, sized in Optimizations. |
| cost | A bulk call that spends 10 tokens is denied unless every bucket holds 10, and a lease of g requests takes g × 10. Rules reject a burst smaller than the largest cost (422). |
Data model.
Three places hold state. Bucket state is hot and disposable; rules are small, audited and versioned; each pod keeps a working copy of both in memory.
Sub-millisecond atomic scripts, and the state is disposable.
| Column | Type | Key | Note |
|---|---|---|---|
| key | text | primary | hash-tagged so an account's buckets share a slot |
| tokens_milli | int | integer milli-tokens, no float drift | |
| ts_ms | int | last refill time from Redis TIME |
| Column | Type | Key | Note |
|---|---|---|---|
| key | text | primary | |
| member | text | one per accepted attempt, at most 5 per IP or 20 per username | |
| score | int | attempt time in ms from Redis TIME; older than the window is trimmed first |
| Column | Type | Key | Note |
|---|---|---|---|
| key | text | primary | |
| count | int |
Small, relational and audited; every change is a new version.
| Column | Type | Key | Note |
|---|---|---|---|
| rule_id | text | primary | |
| match_plan | text | ||
| match_route | text | ||
| identity | enum | api_key | ip | user | |
| algorithm | enum | token_bucket | sliding_log | fixed_window | |
| rate_per_s | int nullable | token_bucket | |
| burst | int nullable | token_bucket | |
| limit | int nullable | count per window: sliding_log, fixed_window | |
| window_s | int nullable | ||
| mode | enum | enforce | shadow | |
| fail_mode | enum | open | closed | |
| version | int |
| Column | Type | Key | Note |
|---|---|---|---|
| account_id | text | primary | |
| rule_id | text | primary | |
| rate_per_s | int | ||
| burst | int | ||
| expires_at | timestamptz nullable | ||
| reason | text |
| Column | Type | Key | Note |
|---|---|---|---|
| rule_id | text | primary | |
| limit_key | text | primary | |
| region | text | primary | |
| share | float | shares of one key sum to 1.0 | |
| updated_at | timestamptz |
Read on every request; rebuilt from the rules store and the counter store within seconds after a restart.
| Column | Type | Key | Note |
|---|---|---|---|
| rule_id | text | primary | |
| version | int |
| Column | Type | Key | Note |
|---|---|---|---|
| limit_key | text | primary | |
| tokens_left | int | ||
| expires_at_ms | int |
| Column | Type | Key | Note |
|---|---|---|---|
| limit_key | text | primary | |
| denied_until_ms | int | now + retry_ms from the script |
| Column | Type | Key | Note |
|---|---|---|---|
| limit_key | text | primary | |
| tokens | int | limit ÷ 50 × 2, used only while the shard is unreachable | |
| ts_ms | int |
| Query | Uses | How |
|---|---|---|
| Is this request allowed? | bucket | one EVALSHA on the account's shard, unless a lease or the deny cache answers first |
| Which rules apply to plan X, route Y? | rule_cache | in memory, per request |
| What changed since version v? | rules | watch from the pod's last version |
| What share of this key's limit does my region get? | region_budgets | watched with the rules, applied to rate and burst |
Interface.
Most deployments run the check as a library inside the gateway. When it runs as a sidecar or a separate service, the call below is the contract. Its shape follows Envoy's rate limit service: a domain plus a list of entries, and the overall verdict is deny if any entry is over. The 429 and its headers are covered in Headers, 429s and backoff.
Checks and spends every entry for one request in one atomic step. The verdict is the strictest entry's.
Admin endpoints
| Method | Path | Does |
|---|---|---|
| PUT | /v1/rules/{id} | Replace a rule. Needs If-Match with the current version. |
| POST | /v1/overrides | Raise or lower one account's limit, usually with expires_at (a merchant's launch week). |
| POST | /v1/rules/{id}:shadow | Put a rule in shadow mode: evaluated and counted, never enforced. |
| GET | /v1/keys/{key}/usage | Tokens left, recent rate and region shares for one key, for support. |
Errors
| HTTP | Type | Body code | Client behaviour |
|---|---|---|---|
| 409 | retry | version_conflict | The rule changed since the caller read it. Re-read and re-apply. |
| 422 | error | invalid_rule | Rejected before it can hurt: rate ≤ 0, burst below the largest cost, or a match that covers every route. |
| 504 | network | store_timeout | The shard took over 5 ms. Apply each rule's fail_mode: open allows and counts fail_open; closed denies with 503. |
| 503 | network | store_unavailable | The shard's breaker is open. Open-mode rules use the pod's fallback bucket at limit ÷ 50 × 2; closed-mode rules (login) deny. |
Optimizations.
The design so far costs one store round trip per request. At Shiplane's traffic mix these changes cut that by about 4×, protect the store from floods and hot keys, and extend the limit across regions.
Leases on a hot key
- A marketplace partner's traffic
- 20K req/son a 25,000/s override, burst 50,000 (assumption)
- Pods it lands on
- 50
- Lease size granted
- 50min(50, 50,000 ÷ 100) = 50 while the bucket is full
- Rate per pod20K ÷ 50400 req/sfrom A marketplace partner's traffic and Pods it lands on
- Store calls per pod400 ÷ 508/sfrom Rate per pod and Lease size granted
- Store calls for this key8 × 50400/s (was 20,000/s)from Store calls per pod and Pods it lands on
- Most tokens out on leases50 pods × 502,500 (5% of the burst)from Pods it lands on and Lease size granted · Never more than half of what the bucket holds, by construction of the grant.
- One hot key is one slot on one primary; leases turn 20,000 calls a second there into 400.
- Twenty-five pods holding leases of 20 can strand at most 500 tokens, which is why leases shrink as the bucket empties.
- Hot counters (sharded-counters) covers contended counters in general; its sharded, slightly stale sums do not suit a limiter, which must decide on a current count, so the limiter leases instead.
Bucket updates per second in us-east
Data
| Stage | updates/s |
|---|---|
| Per-request checks | 360,000 (100%) |
| + deny cache | 324,000 (90%) |
| + leases | 81,000 (23%) |
One limit, three regions
A round trip from us-east to ap-south takes on the order of 200 ms, so a single global counter breaks requirement #2 by two orders of magnitude. Instead, each region gets a share of each key's limit, set by where the key's traffic actually is.
| Step | us-east | eu-west | ap-south |
|---|---|---|---|
| Default split (60/25/15) | 600/s | 250/s | 150/s |
| Observed usage, last 10 s | 70% | 30% | 0% |
| Raise any region under 10% to 10% | – | – | 10% |
| Split the other 90% as observed (70:30) | 63% | 27% | 10% |
| New budgets for a 1,000/s key | 630/s | 270/s | 100/s |
- 01One gateway
- Load
- 2K req/s
- Bottleneck
- None yet
- Change
- In-process token buckets on the one gateway
- Adds
- 1API clients2Gateway pods7Shiplane API
- 0210 pods
- Load
- 20K req/s
- Bottleneck
- Limit ÷ 10 per pod; pinned connections get a fraction of their limit
- Change
- One shared Redis with an atomic script; rules move out of config files
- Adds
- 3Counter store4Rules API5Rules store
- 0350 pods
- Load
- 180K req/s per region
- Bottleneck
- One Redis at 360K bucket updates/s
- Change
- Redis Cluster with 16 primaries, keys hash-tagged per account
- Adds
- 8Metrics
- 04Hot keys
- Load
- One key at 20K req/s
- Bottleneck
- One slot, one primary; store cost per request
- Change
- Local buckets, leases and the deny cache in every pod
- 053 regions
- Load
- 300K req/s
- Bottleneck
- Cross-region round trips of 70 to 200 ms per check
- Change
- A counter store per region, each key's limit split by traffic share
- Adds
- 6Budget rebalancer
Trade-offs.
The chosen option is first; the others stay visible so the reasoning can be checked.
- Pro:One count per key, whichever pod the request hits
- Pro:Sub-millisecond atomic scripts
- Pro:Easy to reason about and to size
- Con:A network hop on the request path
- Con:A store to run, and to plan for when it fails
Wrong whenever traffic is uneven (40/s instead of 1,000/s); Limits drift with autoscaling
A hot key pins one pod; Scaling or a crash moves keys and loses their state
Convergence takes tens to hundreds of ms and bursts slip through meanwhile; Traffic between pods grows with pods × keys
- Pro:An outage of the limiter never becomes an outage of the API
- Pro:Still bounded: at most 2× the limit while the store is gone
- Con:A client can briefly exceed its paid limit
- Con:Needs alerting on the fail_open rate so it is not silent
Every covered request fails while the store is down
Real systems default to open. Stripe's limiters catch any error in the limiter code or Redis and let the request through. Envoy's rate limit filter allows the request when the rate limit service fails, and counts it as failure_mode_allowed; setting failure_mode_deny returns a 500 instead. Shiplane chooses per rule, with fail_mode on each.
A pod's view of one counter shard
States of2Gateway pods
Each pod keeps one of these per shard. It stops paying the 5 ms timeout on every request once a shard is clearly gone, and lets one probe a second find out when it is back.
| From → To | Event | Action | Actor |
|---|---|---|---|
| Global checks → Global checks | reply under 5 ms | reset timeout count | |
| Global checks → Fallback | 3rd timeout in a row | start fallback, fail_open += 1 | |
| Fallback → Probing | 1 s timer | pod timer | |
| Probing → Global checks | probe replies | drop fallback buckets | |
| Probing → Fallback | probe times out |
- Global checksstart
- every check goes to the shard
- Fallbackerror
- open rules use pod buckets; closed rules deny
- Probing
- one real check goes to the shard
Shard 7 fails over
Timeline as a list
Shard 7 fails over: 4 lanes, from 0 s to 24 s.
- 0–3 s · Shard 7 primary · serving
- 0–19 s · Shard 7 replica · async copy
- 0–3 s · Gateway pods · global checks
- 3–24 s · Shard 7 primary · down
- 3–20 s · Gateway pods · fallback: limit ÷ 50 × 2 (allowed)
- 3–20 s · Login rule on shard 7 · fail closed (503) (denied)
- 3 s · all lanes · primary dies (error)
- 3 s · Gateway pods · breaker opens
- 10 s · Gateway pods · probe times out (tick, denied)
- 18 s · all lanes · marked FAIL (15 s) (deadline)
- 19–24 s · Shard 7 replica · primary
- 19 s · Shard 7 replica · promoted (ok)
- 20–24 s · Gateway pods · global checks
- 20 s · Gateway pods · probe ok (ok, ok)
- Pro:Every check stays inside its region
- Pro:Budgets sum to the limit so there is no overshoot
- Con:A sudden shift of traffic to another region is under-served for up to 10 s
- Con:A rebalancer and its inputs to run
A client spread over 3 regions gets up to 3× its limit
70 to 200 ms added to requests far from it; One region's outage breaks limiting everywhere
- Pro:The limit holds to within the lease error
- Pro:Denials are immediate
- Con:A store call (or lease) on the request path
A flood gets through until the count catches up; Not acceptable for quotas a customer pays for
| Failure | Impact | Detection | Mitigation | Meanwhile |
|---|---|---|---|---|
| A shard's primary dies3Counter store | Checks on its 1/16 of keys time out | Pod breakers open; fail_open rate alert | Replica promoted after the node timeout; pods probe once a second; cluster-require-full-coverage no, so the other 15 shards keep serving | API rules run on pod buckets (up to 2× limit); login denies |
| Network partition between pods and the store2Gateway pods | Every check waits for the 5 ms timeout until breakers open | Store latency and timeout rate per pod | Breaker per shard stops paying the timeout | Each rule's fail_mode applies |
| Pod clocks drift2Gateway pods | None on buckets | NTP offset metric | The script reads Redis TIME; pod clocks never enter the refill | Lease expiry (1 s, pod-local) is off by the drift at most |
| A bad rule is pushed (rate 0 on every route)4Rules API | Everyone is denied | Shadow-mode deny rate before enforcing; deny rate by rule | Validation (422), then shadow, then canary pods first; pods keep the last good version | Pods ignore a version that fails their own checks |
| One hot key saturates its primary3Counter store | Latency for every key on that primary | Per-slot ops and the top keys by calls | Leases and deny cache; if still hot, migrate its slot to a dedicated primary | Other keys on the shard see higher p99 until moved |