Skip to content

The finished design, decision by decision

How to design an API Rate Limiter

Not the one correct diagram, but a design you can defend under these constraints: the finished architecture, then every stage's question with the reasoning that answers it, the tradeoffs it accepts, and where another engineer could land differently.

Rate limiting a public API

The short answer

6 parts, each with one job. The map below shows how requests and data move between them; the stages after it explain why each part is there.

API clients
Customer integrations, scripts and dashboards.
Load balancer
Spreads requests across instances; no custom limiting.
API instances
Authenticate the key, check the limit with one atomic Redis call, then handle or reject with 429 and Retry-After. Fall back to local buckets if Redis is unavailable.
Redis
Token buckets per key, updated by an atomic script; expire when idle.
Postgres
API keys and plans; the resource the limiter protects.
Search cluster
Expensive queries whose cost varies by three orders of magnitude.
12345CLIENTAPI clientsEDGELoad balancerSERVICEAPI instancesCACHERedisDATABASEPostgresSERVICESearch cluster

Select a component to see what it is responsible for and which state it owns.

  1. 1API clients → Load balancer: Requests with API key
  2. 2Load balancer → API instances: Round-robin across instances
  3. 3API instances → Redis: Atomic take-tokens script
  4. 4API instances → Postgres: Admitted requests; plan lookups (cached)
  5. 5API instances → Search cluster: Admitted searches, cost-weighted

Why does this design work?

Each key gets one token bucket in a shared store, updated by one atomic script that uses one clock. That gives every instance the same answer without coordinating with each other: the coordination happens in a single, fast, atomic step next to the data.

The rest of the design is about where that ideal bends. Under Redis failure the limiter degrades to local approximation rather than failing the API. For hot keys it batches or splits coordination. For adversaries it counts by several keys at once. For expensive endpoints it counts cost and concurrency rather than requests. Each bend trades precision for availability, throughput or relevance, deliberately and visibly to clients through headers.

Invariants, and where they are enforced

  • A key's admitted rate stays near its plan regardless of how many instances serve it.

    One shared bucket per key, updated by a single atomic script that reads, refills and takes tokens in one step using Redis's clock. Enforced by Redis, API instances.

  • A Redis failure degrades limiting precision, never API availability.

    A short timeout on the limit check, then local per-instance buckets sized to the plan divided by the instance count. Enforced by API instances.

  • Every rejection tells the client when it may retry.

    429 responses carry Retry-After computed from the bucket's deficit and refill rate. Enforced by API instances.

What does it rely on?

  • Redis script execution is atomic per shard, and its round-trip latency stays well under the 2 ms budget.
  • Losing bucket state is acceptable: buckets restart full.
  • Clients honour Retry-After well enough that rejections do not turn into retry storms.
  • Plan changes reach instances within seconds via a short-lived local cache.

What tradeoffs does it make?

ChoiceGainsCosts
Shared Redis bucket per keyFleet-wide correctness with one round trip.A dependency on every request's critical path.
Fail open to local limitsThe limiter never causes an outage.Approximate limits during failures.
Token leases for hot keysTwo orders of magnitude fewer Redis calls.Stranded tokens and small over- or under-admission.
Cost-weighted tokensLimits track the scarce resource.Cost must be estimated before work is done.

What are the reasonable alternatives?

API gateway with built-in limiting (Envoy, Kong, a cloud gateway)
Better when your load balancer or gateway layer supports it: one less component to build, at the cost of flexibility in keys and costs.
Sticky routing with local buckets
Better when A gateway can route consistently by key and the fleet is stable, which removes the network hop entirely.
Approximate distributed counting (gossiped local counts)
Better when global volume is extreme and a few percent of over-admission is fine everywhere.

When does it stop working?

  • The latency budget shrinks below a network round trip, which forces local or gateway-level limiting.
  • Limits must be exact for billing, which calls for metering after the fact rather than limiting.
  • Abuse comes from many keys that look individually normal, which needs anomaly detection rather than limits.
  • Multi-region deployment requires global limits, which bring cross-region coordination or explicit per-region quotas.

Every stage, decided and explained

Spoilers, for the whole investigation: each stage's question and its answer, the reasoning behind it, and the tradeoffs it accepts. If you have not worked through the stages yet, you may want to do that first.

Work through the stages

Stage 1 of 10 · Model

What the numbers say

Thirty instances, 40,000 requests a second, plans measured per minute. Before choosing a design, check which simple ideas the arithmetic already rules out.

What you need to know first

A rate limit is a promise per key: "this API key gets 600 requests a minute". The hard part is that the key's requests don't arrive at one place. The load balancer spreads them round-robin over every instance, so with 30 instances each one sees about 1/30th of any key's traffic.

A Pro key sends 600 requests a minute, spread evenly over 30 instances. About how many of them does one instance see per minute?

About 20 per minute.

600 ÷ 30 = 20 a minute per instance.

So an instance that counts only its own traffic sees 20 requests from a key that is already at its limit. It has no way to tell, from what it sees, that the key is at 600 fleet-wide.

The simplest limiter counts requests per key in fixed calendar windows: "requests in 12:00:00–12:00:59". At the start of each minute the counter resets to zero.

That counts an average over the window. It says nothing about how the requests are spread inside the window, or across the boundary between two windows.

Limit: 60 per calendar minute. A client sends 60 requests at 12:00:59 and 60 more at 12:01:00. What does a fixed-window counter do?

It admits all 120. The first 60 land in the 12:00 window, which still had room; the next 60 land in a fresh 12:01 window that just reset to zero.

120 requests in about one second, from a client "limited" to 60 a minute.

Two numbers to keep in mind for anything on the request path:

OperationRough cost
Redis round trip in the same region~0.3 ms
Simple Redis operations, one shard~100,000 per second
This API's latency budget for the limiter2 ms at p99

A check that touches shared state on every request is affordable here. The question is what happens to every request when that shared state is slow or gone.

Every request now makes one Redis call. Which concern is the real one?

Redis is now on the path of every request, so its latency and failures become the API's.

Before, a Redis problem was a Redis problem. Now every request waits for it. That is the design question the rest of this investigation keeps returning to.

What the stage asks

Which statements hold?

  1. Holds

    If each of 30 instances enforces 600 requests/minute per key in its own memory, a Pro customer can make up to 18,000 requests a minute.

    The load balancer spreads the customer's requests over every instance, and each instance gives them a full allowance. The effective limit scales with the fleet, and autoscaling makes it move.

  2. Holds

    A fixed one-minute window with a limit of 60 can admit 120 requests within two seconds.

    60 in the last second of one window, 60 in the first second of the next. Fixed windows enforce an average over the window, not a maximum over any interval.

  3. Depends

    Checking every request against Redis means about 40,000 Redis operations a second, which is a problem.

    A single Redis shard handles on the order of 100,000 simple operations a second, and a cluster spreads keys over shards, so throughput is fine. What matters is that Redis is now on the critical path of every request: its latency and availability become the API's. That is the real design question.

  4. Fails

    Limiting by client IP is sufficient for the authenticated API.

    One corporate NAT can hide hundreds of honest developers behind one IP, while one customer can spread a script across many machines. The plan belongs to the API key, so the limit must too. IP limits are for traffic with no better identity, such as login.

  5. Depends

    The limiter must be exact: admitting the 601st request in a minute is a bug.

    Plans are commercial promises, and protecting the database is an engineering goal. Neither is harmed by a few percent of imprecision. Insisting on exactness costs coordination and latency on every request; a limiter that is approximately right and always fast is usually the better product.

The reasoning

  1. Per-instance limits multiply by the number of instances, so a fleet-wide limit needs shared state.
  2. Fixed windows enforce an average per window and allow double bursts at the boundary.
  3. Shared state on every request puts its latency and availability on the API's critical path.

Two conclusions shape everything after this:

  • The state must be shared. Per-instance limits multiply by the fleet size, so a key's counter has to live in one place every instance consults, or be split deliberately.
  • The shared state is now on the hot path. Every request waits for it, so its latency budget is the 2 ms requirement and its failure mode is the API's failure mode, unless you design otherwise.

And one framing: rate limiting is a Backpressure and capacity mechanism. It turns "more load than the database can take" into "a clear signal to the client that sent it".

Stage 2 of 10 · Decide

Choose the algorithm

Good clients are bursty: a dashboard loads and fires 20 requests at once, then nothing for a minute. Bad clients are sustained: a loop that never stops. You want to allow the first and stop the second, with little memory per key across ~50,000 keys.

What you need to know first

There are four common rate-limiting algorithms. They differ in what they remember per key:

AlgorithmState per keyEnforces
Fixed windowone counterN per calendar window
Sliding logevery accepted timestampN in any window-length span
Sliding window countertwo countersan estimate of N in the last window
Token buckettokens and a timestampa burst of B, then a steady rate r

Try each pattern below. Watch the most in any 10 s figure for each algorithm: that is how much traffic the limiter really lets through in a short span, whatever the plan says.

Same requests, four limiters. The limit is 5 requests per 10 seconds. Pick a pattern or click the timeline to add requests.
Requests (10)
Fixed window10 through · most in any 10 s: 10
Sliding log5 through · most in any 10 s: 5
Sliding window counter6 through · most in any 10 s: 6
Token bucket6 through · most in any 10 s: 6
15.0s
Fixed window:
Counts per calendar window; the count resets at 10 s.
Sliding log:
Keeps every accepted timestamp from the last 10 s.
Sliding window counter:
Weights the previous window's count by how much of it still overlaps: an estimate, slightly off either way.
Token bucket:
Holds 5 tokens and refills one every 2 s, so a full bucket plus refills can admit a little over 5 in a 10 s span.

With the burst at the boundary, which algorithm lets through twice the limit in under four seconds?

Fixed window

Five requests fall at the end of the first window and five at the start of the second. Each window sees only five, so all ten pass.

A token bucket has two parameters. The capacity B is how big a burst the key may send at once. The refill rate r is the sustained rate. Each request takes a token; tokens refill continuously at r up to B.

When the bucket is empty, the time until the next token is (1 − tokens) ÷ r. That number is what a rejection can send back as Retry-After.

A sliding log stores one 8-byte timestamp per accepted request. An Enterprise key allowed 100,000 requests a minute keeps how many kilobytes of timestamps?

About 800 KB.

100,000 × 8 bytes = 800,000 bytes ≈ 800 KB for one key, rewritten as it slides.

A token bucket for the same key is two numbers, about 16 bytes. Sliding logs are exact, but their memory grows with the limit, and the biggest limits belong to the busiest keys.

An honest dashboard fires 20 requests on page load, then nothing for a minute. The plan is 600 a minute. Under a token bucket with r = 10/s, what decides whether the dashboard gets 429s?

The capacity B: if B is at least 20, the burst fits.

A full bucket absorbs a burst up to B at once. The rate r only matters once the bucket is empty.

What the stage asks

Which algorithm should enforce plan limits?

  1. Defensible

    Fixed window: count requests per key per calendar minute

    One counter per key and trivially cheap, which is why many systems start here. But it allows double-rate bursts at window boundaries, and every client's allowance resets at the same instant, which synchronizes retrying clients into a spike at the top of each minute.

  2. Defensible

    Sliding log: store every request's timestamp and count those in the last 60 s

    Exact over any interval, but memory grows with the limit: an Enterprise key at 100,000 a minute keeps 100,000 timestamps. Precise where precision is not required, and expensive where you have the most traffic.

  3. Sound

    Token bucket: capacity B for bursts, refilling at the plan's rate

    Two numbers per key (tokens and last-refill time) express exactly what the requirements say: allow a burst up to B, then hold the client to the sustained rate r. A dashboard's 20-request burst fits in the bucket; a loop drains it and then gets r per second. Retry-After falls out of the arithmetic: the deficit divided by r.

  4. Defensible

    Limit concurrent in-flight requests per key instead of rate

    Concurrency limits are excellent at protecting a backend from slow, expensive work, and you will want one later. But plans are sold as requests per minute, and a client making fast requests one at a time can still exceed its plan many times over.

What a strong answer covers

  • Separates burst allowance from sustained rate, so short bursts pass and long ones are held to the plan.
  • Considers memory per key: constant state vs. state that grows with the limit.
  • Recognizes the fixed window's boundary effect or synchronized resets.Supporting

The reasoning

  1. A token bucket separates the burst a key may send (capacity B) from the rate it may sustain (refill r).
  2. Sliding logs are exact but store state proportional to the limit; buckets store two numbers.
  3. Fixed windows are cheap but admit double bursts at boundaries and reset everyone at the same instant.

A token bucket encodes a contract: "you may burst up to B, and sustain r". For Pro, r = 10/s and B = 50 might be right: a dashboard's burst passes, while a loop gets 10 a second, which is 600 a minute as sold.

B is a product decision dressed up as a parameter. Too small and well-behaved clients get 429s on page load; too large and an abusive client gets a big free burst every time it pauses. Pick it from observed honest bursts, not from the plan.

Tradeoffs

ChoiceGainsCosts
Token bucketConstant memory; bursts and sustained rate both expressed.Two parameters to tune; harder to explain than 'N per minute'.

Stage 3 of 10 · Decide

Where does the count live?

Thirty instances, autoscaling between 20 and 60, round-robin load balancing. Every instance must agree, closely enough, on how many tokens each key has left.

What you need to know first

Wherever the bucket lives, three properties matter:

  • Shared: every instance must consult the same bucket for a key, or the limit multiplies by the number of instances.
  • Fast: the check runs on every request, inside a 2 ms budget.
  • Atomic: many instances update the same key at the same moment, so the read-refill-take-write sequence must not interleave.

Durability matters much less. If bucket state is lost, buckets start full, which briefly admits a little extra traffic.

The store holding every token bucket loses all its data. What happens?

Every key gets a full bucket again: a short burst of extra traffic, then normal limiting.

A missing bucket is treated as a full one. The worst case is one extra burst of B per key. That is why rate-limit state is a good fit for an in-memory store without strong durability.

Consider where else the count could live:

  • Each instance's memory costs nothing, but it isn't shared.
  • Sticky routing, sending every request for a key to one instance, makes local memory shared again. It needs a load balancer that can route by key, and every scale event moves keys between instances.
  • The primary database is shared and durable, but every request would write to the very database the limiter exists to protect.
  • Redis is shared, about 0.3 ms away, and can run a small script atomically.

If the bucket were a Postgres row updated in a transaction, what happens when one key receives 500 requests a second from 30 instances?

500 transactions a second all try to lock the same row. They queue behind each other on the row lock, latency climbs, and the primary database (the thing being protected) spends its capacity on limiter writes.

The limiter turns a traffic spike into exactly the database overload it was supposed to prevent.

What the stage asks

Where should token buckets live?

  1. Flawed

    In each instance's memory

    Zero latency and the wrong answer: as the first stage showed, the effective limit becomes the plan times the instance count, and it changes as the fleet scales.

  2. Defensible

    Route each key to one instance (consistent hashing) and keep its bucket there

    Correct while the fleet is stable, with no network hop. But the load balancer cannot do it, every scale event reshuffles keys and resets their buckets, and the largest customer's traffic all lands on one instance. It is a good design for a gateway built for it, and an awkward one here.

  3. Sound

    In Redis, updated by one atomic script per request

    One bucket per key, shared by every instance, updated in a single round trip (~0.3 ms). The script reads, refills, takes and writes as one atomic step, so concurrent requests from thirty instances cannot interleave. Redis is now on the critical path, which the failure stage deals with.

  4. Flawed

    In a Postgres row per key, updated in a transaction

    Every request now writes to the very database you are protecting, and the hottest keys become lock-contended rows. The limiter becomes the overload it was meant to prevent.

What a strong answer covers

  • All instances must see one counter per key, or the limit multiplies.
  • Updates must be atomic, because many instances modify the same key concurrently.
  • Notes the new dependency on the request path and its latency budget.Supporting

The reasoning

  1. Rate-limit state must be shared across instances and updated atomically, but it doesn't need to be durable.
  2. Losing a bucket means it starts full: a short burst of extra admissions, not an outage.
  3. Keeping counters in the database you are protecting turns the limiter into a source of load.

Rate-limit state is unusually forgiving: losing it means buckets start full, which briefly admits a little extra traffic. That makes an in-memory store like Redis a good fit: fast, shared, and without durability that nobody needs. See Soft state.

What is not forgiving is concurrency. Thirty instances touching the same key at the same moment is the normal case, and the next stage is what happens when the update is not atomic.

Stage 4 of 10 · Break it

The limiter that leaks

The buckets are in Redis, as decided. The logic is a direct translation of the token bucket. Find the lines that let a 60/minute key make 4,000 requests, and the line that made the incident worse.

What you need to know first

A read-modify-write reads a value, computes a new one, and writes it back. If it takes two round trips (a GET, then a SET), other clients can read the same old value in between.

Each of them computes its own update from the same starting point, and the last write wins. The other updates are lost.

A bucket holds 1 token. Three instances each GET it at the same moment, see 1 token, allow their request and SET the bucket to 0. How many requests were admitted, and how many tokens were spent?

Three admitted, one token spent.

All three acted on the same stale read. Three SETs of 0 leave the bucket at 0, as if one request had run. Across 30 instances this is how a 60-a-minute key makes thousands of requests.

Two more details decide whether a limiter behaves:

  • Whose clock? Refilling needs "seconds since the last refill". If each instance uses its own clock, one running ahead refills tokens that were never earned. One clock for every caller fixes that, for example the store's own clock.
  • What does a rejection say? Most HTTP clients retry failed requests. A 429 with no hint about when to come back is usually retried at once, so every rejection produces another request.

A client retries every failed request immediately, and its limiter rejects 90% of its requests. What happens to the number of requests the limiter has to handle?

It multiplies. Each rejected request comes straight back, and most of those are rejected again. A client that wanted 100 requests can generate around 1,000 attempts.

A Retry-After header tells well-behaved clients how long to wait, which turns a retry storm into a pause.

What the stage asks

Select the faulty lines.

Codetypescriptlimit.ts
  1. 1async function allow(key: string, plan: Plan): Promise<boolean> {
  2. 2 const raw = await redis.get(`rl:${key}`);
  3. 3 const state = raw ? JSON.parse(raw) : { tokens: plan.burst, ts: Date.now() };
  4. 4 const elapsed = (Date.now() - state.ts) / 1000;

    Each instance uses its own clock. An instance whose clock runs ahead computes extra elapsed time and refills tokens that do not exist. Use one clock: Redis's TIME, inside the script.

  5. 5 const tokens = Math.min(plan.burst, state.tokens + elapsed * plan.perSecond);
  6. 6 if (tokens < 1) return false;

    The rejection carries no Retry-After. The customer's client, like most, retried immediately, so every 429 caused another request. Rejections must tell clients when capacity returns.

  7. 7 await redis.set(`rl:${key}`, JSON.stringify({ tokens: tokens - 1, ts: Date.now() }));

    Read-modify-write across two round trips. Thirty instances read the same state, each subtracts one, and each writes back: 30 requests cost one token. The whole update must be one atomic operation (a Lua script), and the key needs a TTL so idle keys disappear.

  8. 8 return true;
  9. 9}

What the fix has to do

  • Concurrent read-modify-write from many instances loses updates, so many requests share one token.
  • The fix is a single atomic operation on the server (script or transaction), not a check in application code.
  • Refill time must come from one clock, not each instance's.
  • 429s without Retry-After invite immediate retries that amplify the load.Supporting

The reasoning

  1. A read-modify-write in two round trips lets concurrent callers act on the same stale value: updates are lost.
  2. Move the whole decision into one atomic operation next to the data, using one clock.
  3. Rejections without Retry-After invite immediate retries that multiply the load.

The token-bucket arithmetic was correct; the concurrency was not. "Read the counter, decide, write the counter" is the check-then-act pattern from every other investigation, and it breaks the same way here: two readers act on the same stale value. See Concurrency control.

The fix moves the whole decision into Redis, where a Lua script runs atomically with respect to other commands on that shard. That makes it one round trip, one clock and one consistent answer.

Stage 5 of 10 · Break it

Write the atomic bucket

Write the Redis-side logic (Lua, or pseudo-code for "runs atomically on the server") that refills, takes tokens and reports the result, plus the caller that sets the response headers.

What you need to know first

Redis runs a Lua script as one atomic step: no other command on that shard runs while the script is running. So a script can read a bucket, refill it, take tokens and write it back without anyone else seeing it half-updated.

Inside a script, redis.call('TIME') returns the server's clock as seconds and microseconds. Every instance that calls the script shares that one clock.

The refill arithmetic, for a bucket with capacity capacity and refill rate tokens per second:

tokens = min(capacity, tokens + (now − last) × rate)

If tokens ≥ cost, subtract the cost and allow. Otherwise reject, and the wait until enough tokens exist is (cost − tokens) ÷ rate.

A bucket refills at 10 tokens a second and holds 0.4 tokens. A request costs 1. How many seconds until it can be admitted?

About 0.06 seconds.

(1 − 0.4) ÷ 10 = 0.06 seconds. As a Retry-After header, round up to whole seconds: 1.

Why give each bucket key an expiry (TTL)?

So idle keys disappear; a missing bucket is just a full one.

After capacity ÷ rate seconds without requests, a bucket is full again. Deleting it then loses nothing and keeps memory proportional to active keys.

Redis converts Lua numbers to integers when a script returns them. A bucket holds 4.7 tokens. What should the script return?

The value as a string, parsed by the caller.

Strings pass through unchanged, so fractional tokens and wait times survive. Returning the raw number would truncate 4.7 to 4 and 0.06 seconds to 0.

What the stage asks

Implement an atomic token bucket and its caller.

Reference implementation

-- token_bucket.lua
local capacity = tonumber(ARGV[1])
local rate     = tonumber(ARGV[2])
local cost     = tonumber(ARGV[3])
local t   = redis.call('TIME')
local now = tonumber(t[1]) + tonumber(t[2]) / 1e6

local b = redis.call('HMGET', KEYS[1], 'tokens', 'ts')
local tokens = tonumber(b[1]) or capacity
local ts     = tonumber(b[2]) or now
tokens = math.min(capacity, tokens + math.max(0, now - ts) * rate)

local allowed = 0
local wait = 0
if tokens >= cost then
  tokens = tokens - cost
  allowed = 1
else
  wait = (cost - tokens) / rate
end

redis.call('HSET', KEYS[1], 'tokens', tokens, 'ts', now)
redis.call('EXPIRE', KEYS[1], math.ceil(capacity / rate) + 1)
return { allowed, tostring(tokens), tostring(wait) }

// limit.ts
async function checkLimit(key: string, plan: Plan, cost = 1) {
  const [allowed, tokens, wait] = await redis.evalsha(
    SCRIPT_SHA, 1, `rl:{${key}}`, plan.burst, plan.perSecond, cost);
  return {
    allowed: allowed === 1,
    headers: {
      "X-RateLimit-Limit": String(plan.perMinute),
      "X-RateLimit-Remaining": String(Math.floor(Number(tokens))),
      ...(allowed === 1 ? {} : { "Retry-After": String(Math.ceil(Number(wait))) }),
    },
  };
}
  • Redis runs a script without interleaving other commands on that shard, so concurrent callers are serialized. This is the atomicity the leaky version lacked.
  • TIME is read inside the script: every instance shares one clock.
  • The TTL is the time to refill completely. After that, a missing key and a full bucket are the same thing, so expiry loses nothing. This is Soft state at work.
  • Numbers are returned as strings because Redis truncates Lua numbers to integers on the way out.

What a strong answer covers

  • Read, refill, take and write happen in one server-side atomic script.
  • Time comes from the server (Redis TIME), not the caller.
  • Refill is capped at capacity, and tokens are only taken when enough are available.
  • Returns the time until enough tokens exist, used for Retry-After.
  • Sets a TTL so idle buckets expire (a missing bucket is simply full).Supporting
  • Supports a cost per request rather than assuming 1.Supporting

The reasoning

  1. A Lua script runs atomically on its Redis shard, so read, refill, take and write happen as one step.
  2. Read the time inside the script (Redis TIME) so every caller shares one clock.
  3. Expire idle buckets after they would have refilled; a missing bucket is a full bucket.

The script is short because the hard part is not the arithmetic. It is putting the arithmetic where it can run atomically, next to the data, with one clock. Most distributed counters, quotas and inventory checks end up with the same shape.

Stage 6 of 10 · Break it

Redis goes down

The limiter now sits on every request. The requirement is explicit: if its state store fails, the API stays up.

What you need to know first

When a dependency on the request path fails, every request has to do something. There are two broad choices:

  • Fail closed: treat "I can't check" as "no". Safe for the thing being protected, but the dependency's outage becomes your outage.
  • Fail open: treat "I can't check" as "yes". The service stays up, but the protection is gone for as long as the dependency is.

The right choice depends on what the check protects. A payment fraud check might fail closed. A limiter whose job is availability usually shouldn't.

A dependency that is slow is often worse than one that is down. A down Redis returns errors in microseconds. A Redis in the middle of a failover may take seconds to answer, and every request waits that long.

A timeout turns slow into failed. It needs to be shorter than the latency the caller can afford to add, here a few milliseconds.

The API handles 40,000 requests a second. Redis stops answering and each limit check waits for a 2-second client timeout. About how many requests are stuck waiting at once?

About 80,000 requests.

Requests in flight = arrival rate × time each one waits = 40,000 × 2 = 80,000.

That is far more concurrent requests than the instances' thread pools or connection limits can hold. The API stops answering everything, including requests that never needed Redis.

Redis is unavailable, and each of 30 instances must limit a Pro key (600 a minute) using only its own memory. What limit should each instance apply so the fleet stays near the plan?

About 600 ÷ 30 = 20 a minute per instance. The load balancer spreads the key's requests roughly evenly, so 30 local limits of 20 add up to about 600.

It's approximate: if the spread is uneven, some instances reject early while others have room. But the database stays protected and the API stays up.

What the stage asks

What should an instance do when the limit check fails or is slow?

  1. Flawed

    Fail closed: reject requests with 503 until Redis is back

    The limiter becomes a single point of failure for the whole API, which violates the requirement. Every customer is down because a cache failed over.

  2. Sound

    Fail open after a ~5 ms timeout, enforcing local buckets sized to the plan divided by the instance count

    The API stays up and stays protected, approximately: each instance allows about 1/30th of a key's plan, so the fleet-wide total stays near the limit while Redis is unavailable. The short timeout keeps a slow Redis from adding latency to every request. Precision drops for a few minutes; availability does not.

  3. Defensible

    Fail open: skip limiting until Redis recovers

    Availability is preserved, and for a short failover this is often acceptable. But the window is exactly when an abusive client faces no limit at all, and the original incident shows what one client can do to the database in a few minutes.

  4. Flawed

    Retry the Redis call until it succeeds before handling the request

    Every request now waits for the failover. Thread pools fill, timeouts cascade, and the API goes down slowly instead of quickly. A dependency on the hot path needs a strict timeout and a fallback, not patience.

What a strong answer covers

  • The limiter must not take down the thing it protects; failing closed makes it a single point of failure.
  • A strict timeout bounds the latency a slow Redis can add.
  • Some protection should remain during the failure, e.g. local limits of plan / instances.

The reasoning

  1. A limiter must not take down the service it protects: failing closed makes it a single point of failure.
  2. Bound the latency a slow dependency can add with a strict timeout, then fall back.
  3. Degrade precision, not availability: local limits of plan ÷ instances keep rough protection during the outage.

A limiter exists to protect availability, so it must never cost more availability than it saves. The design principle is degrade precision, not service: fall back from a shared, exact-ish limit to a local, approximate one, with a timeout tight enough that the fallback triggers quickly.

Local fallback needs the instance count, which autoscaling changes. Each instance can read the current count from the orchestrator or the load balancer's target list; being off by a few instances only changes the limit by a few percent.

Tradeoffs

ChoiceGainsCosts
Fail open with local bucketsAPI stays up and roughly protected.Limits are approximate, and uneven load across instances can admit somewhat more than the plan.

Stage 7 of 10 · Change it

The biggest customer

Partitioning Redis by key spreads 50,000 keys nicely, except when one key is most of the traffic.

What you need to know first

Redis Cluster splits keys over shards by hashing the key name. 50,000 keys spread evenly, and adding shards spreads them further.

But all commands for one key go to one shard, and each shard runs commands on a single thread. A key that receives most of the traffic stays on one shard however many shards exist. That is a hot key.

One key gets 15,000 script calls a second and its shard is at 90% CPU. What does doubling the number of shards do for that key?

Nothing: the key still lives on exactly one shard.

More shards spread more keys. They can't split one key's traffic, because every call for that key must reach the shard that holds it.

Two ways to take load off one key:

  • Batch: an instance takes many tokens at once (say 200) and spends them locally. One Redis call now covers 200 requests.
  • Split: store the key as N sub-buckets on different shards, each holding 1/N of the limit. Each request picks one sub-bucket.

Both reduce coordination per request. Both make the limit slightly less exact.

15,000 requests a second for one key, with instances taking tokens 200 at a time. About how many Redis calls a second does that key need?

About 75 calls per second.

15,000 ÷ 200 = 75 calls a second, down from 15,000.

The cost: tokens an instance has taken but not yet used are invisible to everyone else. At most 200 per instance are "stranded" at a time, which is about 1% of this key's per-second traffic.

What the stage asks

How should the limiter handle this key?

  1. Defensible

    Move to larger Redis nodes

    It buys time, but one key's commands still run on one shard's single thread. The next customer twice this size hits the same wall.

  2. Sound

    Let instances take tokens in batches (say 200 at a time) and spend them locally

    Redis calls for this key drop by two orders of magnitude, from 15,000 a second to ~75. The cost is precision: tokens leased to an instance that then sees no traffic are briefly stranded. At 15,000 a second, being off by a few hundred is noise.

  3. Sound

    Split the key into N sub-buckets on different shards, each with 1/N of the limit, chosen at random per request

    Spreads the load over shards with no local state. Random choice makes each sub-bucket see roughly 1/N of the traffic, though the variance means the customer may hit one sub-bucket's limit slightly early. Fine for large keys; pointless for small ones.

  4. Flawed

    Exempt Enterprise keys from limiting; they pay for capacity

    The incident that started this came from a paying customer's retry loop. Paying for a higher limit is not the same as paying to remove the database's protection.

What a strong answer covers

  • A single hot key stays on one shard however many shards exist.
  • The fix reduces per-request coordination for that key (batching or splitting).
  • Names the precision lost, and why it is acceptable at this volume.

The reasoning

  1. Partitioning spreads many keys; it can't spread one hot key, which always lives on one shard.
  2. Reduce coordination for hot keys by leasing tokens in batches or splitting the key into sub-buckets.
  3. The precision lost is tiny relative to a very large key's traffic, so the trade gets cheaper as the key grows.

Partitioning scales many keys, not one hot key. For hot keys you split the work: lease tokens in batches so most requests are decided locally, or split the key so the coordination spreads out.

Both trade precision for throughput, and the trade gets cheaper as the key gets bigger: being off by 200 tokens matters for a 60/minute key and is invisible for a 15,000/second one. Hybrid limiters apply the batched path only above a traffic threshold.

Stage 8 of 10 · Change it

Credential stuffing on login

Login is unauthenticated, so there is no API key. The traffic is not one heavy client; it is thousands of light ones.

What you need to know first

Credential stuffing tests leaked email–password pairs from other sites against your login page. Attackers spread the attempts over thousands of IPs, often residential proxies, so each IP makes only a few attempts an hour.

A rate limit catches abuse that is concentrated in the key it counts by. Distributed attacks are built to stay under every per-key limit.

8,000 IPs each try 4 passwords an hour. About how many login attempts is that per day?

About 768,000 attempts per day.

8,000 × 4 × 24 = 768,000 attempts a day, while each IP stays at 4 an hour, below any per-IP limit you could set without blocking an office behind a shared NAT.

So login protection counts by several keys at once:

KeyCatches
Per IPone machine hammering
Per accountmany guesses at one user's password
Per IP + accountone source retrying one user
Global failure ratea distributed attack, visible only in aggregate

And the response can escalate (add a delay, require a CAPTCHA) instead of hard-blocking.

You lock any account for an hour after 5 failed logins. What can an attacker who only knows a victim's email address do?

Lock the victim out on purpose: 5 wrong passwords and the real owner can't log in for an hour. Repeated every hour, it's a denial of service against that user.

That's why per-account limits usually slow down or challenge further attempts rather than blocking the account outright.

What the stage asks

Which statements hold for protecting login?

  1. Fails

    A limit of 5 attempts per IP per minute stops this attack.

    Each attacking IP makes a few attempts an hour, far below any per-IP limit you could set without hurting offices behind one NAT. Distributed attacks are designed to sit under per-source limits.

  2. Holds

    Locking an account after 5 failed attempts lets an attacker lock out real users on purpose.

    Anyone who knows an email address can trigger the lockout. Per-account limits are useful, but their response should slow down or add challenges (delays, CAPTCHA, email verification) rather than block the real owner.

  3. Holds

    A sudden rise in the global failed-login rate is a useful signal even when no single IP or account exceeds its limit.

    Distributed attacks are invisible per key and obvious in aggregate. Global limits and alerts, such as stepping up challenges for everyone when failures spike, catch what per-key limits cannot.

  4. Fails

    With good rate limits, breached-password checks and MFA are unnecessary.

    Rate limiting slows guessing; it does not make a leaked password safe. An attacker with 8,000 IPs and patience still gets through eventually. Limits are one layer of defence.

The reasoning

  1. Distributed attacks stay under per-source limits by design; count by several keys and watch the aggregate.
  2. Hard lockouts let attackers lock real users out; escalate with delays and challenges instead.
  3. Rate limits slow guessing but don't make leaked passwords safe; they are one layer among several.

Rate limits assume that abuse is concentrated in a key. Adversaries deliberately spread it out. Protecting login takes keys at several granularities (IP, account, IP-and-account pair, global), responses that escalate (delay, challenge) rather than hard-block, and signals in aggregate. The same token-bucket mechanism serves all of them; what changes is the key you count by.

Stage 9 of 10 · Change it

Not all requests cost the same

Every customer stayed inside their requests-per-minute plan. The cluster went down anyway.

What you need to know first

A request-count limit assumes requests cost roughly the same. When one request can take 1,000 times longer than another, counting requests says little about load.

What the backend runs out of is time: how many query-seconds it can serve per second. That's the resource to limit.

A customer's plan allows 10 searches a second, and each of their searches takes 2 seconds of cluster time. How many search-seconds of work do they add every second?

About 20 search-seconds.

10 × 2 s = 20 seconds of work per second. On average about 20 of their queries are running at any moment.

Another customer with the same plan whose queries take 2 ms adds 0.02 search-seconds per second. Same request limit, a thousand times the load.

Three tools, each limiting something different:

  • Cost-weighted tokens: a request takes tokens in proportion to its estimated cost, so the bucket measures work rather than requests.
  • Concurrency limit per key: at most N requests from one key running at once, which bounds how much of the backend one client can occupy.
  • Global load shedding: when the backend's queue grows, reject or defer work from everyone, even clients within their own limits.

Every customer is within their own per-key limits, yet together they exceed what the cluster can serve. Which tool still protects it?

A global limit or load shedding based on the cluster's own state

Per-key limits each look fine; only something that watches total load can notice that the sum is too much.

What the stage asks

How should search be protected?

  1. Defensible

    Give search a lower requests-per-minute limit

    It reduces the damage, but a request count cannot tell a 2 ms query from a 2 s one. Cheap queries get throttled unnecessarily, and a few expensive ones can still saturate the cluster.

  2. Sound

    Charge tokens by estimated cost, cap concurrent searches per key, and shed load globally when the cluster's queue grows

    Cost-weighted tokens make the bucket measure the scarce resource rather than the request count. A per-key concurrency cap stops one customer occupying the cluster with slow queries. A global shed protects the cluster when everyone is busy at once. Each layer covers a different failure.

  3. Defensible

    Scale the search cluster until it copes

    Capacity helps the average case. But query cost is unbounded on the customer's side, so a few pathological queries can always outrun capacity. You need a limit on the input, not just more output.

  4. Defensible

    Kill any query that runs longer than 500 ms

    A sensible complement that bounds the worst case per query. On its own, many 450 ms queries still saturate the cluster, and the work done before the timeout is wasted.

What a strong answer covers

  • Limits should measure the scarce resource (cost or time), not just request count.
  • Concurrency limits bound how much of the backend one client can occupy at once.
  • A global limit or load shedding protects the backend when every client is within its own limit.Supporting

The reasoning

  1. When request costs vary widely, limit the scarce resource (cost or time), not the request count.
  2. Per-key concurrency caps stop one client from occupying a slow backend.
  3. A global shed protects the backend when every client is within its own limit at once.

Request counts are a proxy for load, and they stop working when requests differ in cost. The fix is to count what is scarce: tokens weighted by estimated cost (refunded when the actual cost is lower), concurrency per key, and a global load-shedding valve. The bucket script already takes a cost argument; the investment is in estimating cost before running the query.

Stage 10 of 10 · Defend it

Defend a 429

A Pro customer files a ticket: "Our logs show we sent 540 requests in the last minute, under our 600 limit, and we got 429s. Your limiter is broken."

Respond as the engineer who built it: how can that happen, which cases are bugs, and what would you check?

What you need to know first

Customers read plans as "600 requests a minute". A token bucket enforces something slightly different: a burst of up to B at once, then a refill of r per second. Over a long period those agree. Over a few seconds they can differ.

Pro has B = 50 and r = 10 per second. A client sends 200 requests in the first 5 seconds of a minute, then 340 spread over the remaining 55 seconds. Does it get any 429s?

Yes, in the first 5 seconds. The bucket starts with 50 tokens and refills 50 more over 5 seconds, so about 100 of the 200 get through and about 100 are rejected. The rest of the minute is fine: 340 over 55 seconds is about 6 a second, under the refill rate.

The minute's total is 540, under 600. The client is well within its "per minute" plan and still saw 429s, because the bucket limits bursts.

Other reasons a client's count and the limiter's count disagree:

  • The limiter ran in degraded mode (local limits during a Redis failure).
  • The client's HTTP library retried requests that its own logs don't show.
  • Other services share the same API key.
  • The client counts by send time, the server by arrival time.

Each of these is checkable if the limiter records its decisions: per key, whether a request was allowed, the tokens remaining, and whether degraded mode was on.

Which observation would point to a real bug rather than expected behaviour?

429s returned while the X-RateLimit-Remaining header on the previous response was well above zero

The limiter told the client it had tokens left and then rejected it. Either the header or the decision is wrong.

What the stage asks

Explain how a client under its per-minute limit can legitimately receive 429s, which explanations would indicate a bug, and how you would find out which happened.

Reference answer

It is probably not a bug, but you should prove that rather than assert it.

Legitimate causes:

  • Burst, not total. Pro might be "50 burst, 10/s refill". If the customer sends 200 requests in the first five seconds of a minute, the bucket empties after ~100 and later requests are rejected, even though the minute's total ends up at 540. The plan says 600 a minute; the mechanism enforces a burst and a rate.
  • Degraded mode. If Redis was failing over, each instance enforced plan / instance count locally. With uneven distribution, some instances reject while others have spare capacity.
  • Different counts. The customer's logs may miss retries their HTTP library made, or other services using the same key; they count when they sent, and the limiter counts when requests arrived, so a minute's boundary differs.

Bugs would look like: rejections while X-RateLimit-Remaining was well above zero; rejections for a key whose plan was changed but whose cached plan was stale; a clock or script error refilling too slowly.

How to find out: pull the key's limiter decisions for that minute (allowed, tokens remaining, degraded flag), compare with the response headers the customer received, and check Redis health for the period.

What to change: make the contract visible. Document burst and refill rates, return remaining tokens on every response, and send Retry-After on every 429. A limiter customers can predict generates fewer tickets than an exact one they cannot.

What a strong answer covers

  • A token bucket enforces burst capacity and refill rate, not 'N per calendar minute': a burst can exhaust it below the per-minute total.
  • During a Redis failure, local fallback limits are approximate and can reject earlier on unevenly loaded instances.
  • The client's count may differ from the server's: retries, other services sharing the key, or counting by send time vs. arrival time.
  • Proposes evidence: per-key limiter logs or metrics, the remaining-tokens header at each 429, whether degraded mode was active.
  • Suggests making the contract explicit to customers (document burst and refill; expose headers) so behaviour is predictable.Supporting

The reasoning

  1. A token bucket enforces a burst and a rate, not 'N per calendar minute'; make that contract visible to clients.
  2. Record each limiter decision (allowed, tokens left, degraded mode) so disputes are settled with evidence.
  3. A predictable limiter with clear headers produces fewer complaints than an exact one clients can't reason about.

A rate limiter is a contract with clients, and the hardest part of a contract is making it legible. Every mechanism in this investigation (bursts, approximation under failure, token leases for big keys) is a place where "600 per minute" stops being literally true. Defending the design means being able to explain each of those departures, show why it serves the customer, and produce the evidence when asked.

How you did

Now try it as an interview question

  • “Design a rate limiter.”
  • “How would you enforce API quotas across a fleet of stateless servers?”
  • “Design protection against credential stuffing on a login endpoint.”
  • “Your rate limiter's Redis goes down. What happens to your API?”

The interview mode mixes stages from this and other investigations with concept recall and questions about your own projects.

Back to the last stage