System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Rate Limiter: placement, distributed counters and follow-ups


The algorithm is a small part of the design. Where the limiting code runs decides who it protects, what it can see, and who maintains it. And the moment the API runs on more than one server, the naive design is wrong — which is the hard part of the problem, and the reason it is asked.

This lesson picks up from the requirements and algorithms. It places the limiter, then makes its counters correct across a fleet, then designs what the client receives and walks through the follow-ups. What the client receives is part of the design: a limiter that rejects without telling the client how to behave produces retry storms, which is the problem it was meant to prevent.

Option 1: middleware inside the application

Three places the counter can runClient: easily bypassedGateway: protects allSidecar: per-serviceMiddleware: sees user
The gateway is the usual answer because it is the only layer that sees traffic no service ever receives.

The limiter is a function in the request pipeline of each service.

  • For: full access to application context — the authenticated user, the customer tier, which endpoint, how expensive this specific call will be. Nothing new to operate.
  • Against: every service must implement it, in every language you use, and they will drift. The request has already crossed the network and been parsed and authenticated before it is rejected, so a flood still costs you the work of receiving it.
  • Right when: limits are highly application-specific, or the system is one service.

Option 2: at the API gateway or reverse proxy

The limiter runs in the shared entry layer (Load balancing and proxies) that already terminates Transport Layer Security and routes requests.

  • For: one implementation for all services, in a component that already exists. Rejected requests never reach application servers, so the flood is absorbed at the edge. Consistent headers and error responses.
  • Against: the gateway sees the request, not the application's context — it may not know the customer tier without a lookup, and it cannot know a request will be expensive. The gateway becomes a component whose failure affects everything.
  • Right when: you have multiple services and a gateway. This is the recommended default for a public API.

Option 3: a sidecar proxy

A proxy running beside every service instance (the service mesh pattern), applying limits to traffic in and out.

  • For: uniform policy across services without changing any of them, and it works for service-to-service traffic, not only inbound public traffic — which is where cascading overload usually starts.
  • Against: an extra hop each way, and a large operational surface for a team that does not already run a mesh.
  • Right when: you already run a service mesh, or the problem is internal services overwhelming each other.

The honest recommendation

Say this: "Gateway by default, because it protects every service with one implementation and rejects floods before they cost anything. I'd add application middleware for limits that need context the gateway doesn't have — per-tier quotas, or throttling a specific expensive query. A sidecar only if we're already running a mesh."

Two layers is normal, not redundant. A coarse limit at the edge stops volumetric abuse; a fine limit inside enforces business rules.

Why per-server counters undercount

Wherever the limiter runs, it runs on many machines.

Why per-server counters undercount10 servers,10 countersEach allowsthe full limitEffectivelimit is 10xSharedcounter in RedisINCR, so noread firstRead-then-write races; INCR or a Lua script makes it one atomic step.
The naive design is wrong by exactly the number of servers, which is a factor that grows as you scale.

Ten API servers behind a round-robin load balancer, each keeping its own in-memory counter, each enforcing 100 requests per minute per user. A client's requests are spread across all ten, so each server sees roughly 10 requests per minute and allows every one. The client sustains 1,000 requests per minute against a 100 limit — the limiter is off by exactly the server count.

Making it worse: the effective limit changes whenever you scale the fleet, so the same configuration means different things on a Tuesday and a Black Friday.

The fix: shared counters

Move the counters to a store every API server can reach — an in-memory key-value store such as Redis is the standard choice, keyed as ratelimit:{user_id}:{window}.

Cost: one network round trip per request, roughly 0.5 ms in the same datacentre. Requirements and scale established that this is acceptable for a 200 ms API.

The race condition in read-then-write

Here is the naive shared-counter implementation, and it is wrong:

Text
count = store.get(key)          # returns 99if count < 100:    store.set(key, count + 1)   # writes 100    allow()else:    reject()

Two requests for the same user arrive at two different API servers at the same moment. Both read 99. Both conclude 99 < 100. Both write 100. Two requests were allowed where one should have been, and the counter now says 100 when 101 have been served.

At 50,000 requests per second this is not a rare edge case. It happens continuously, and the overshoot grows with concurrency — with 10 servers hammering one hot key, the count can drift far above the limit.

Fix 1: atomic operations

Most in-memory stores provide an atomic increment that returns the new value. One round trip, no read-then-write gap:

Text
count = store.incr(key)          # atomic; returns the new valueif count == 1:    store.expire(key, 60)        # set the window TTL on first useif count > 100:    reject()else:    allow()

This is correct for a fixed window counter and it is the answer to give first. The residual wrinkle: incr and expire are two calls, so a crash between them leaves a key with no expiry. Some stores let you set both in one command; otherwise Fix 2 handles it.

Fix 2: a server-side script

Send a small script that the store executes atomically as one unit — Redis executes Lua scripts this way. The whole decision happens inside the store: read the counter, apply the sliding-window arithmetic from The five algorithms, compared, decide, write, set the expiry, and return the verdict.

  • One round trip regardless of how many operations the algorithm needs.
  • Genuinely atomic, so no other client can interleave.
  • Necessary for token bucket and sliding window algorithms, which need read-compute-write rather than a single increment.

This is the recommended answer for anything beyond a fixed window counter.

A third approach worth naming: for a sliding window log, a sorted set with the timestamp as score, using a transaction that removes expired entries, adds the new one, and counts, all in one atomic block.

Fix 3: local counters with asynchronous synchronisation

When the 0.5 ms round trip is too expensive, each server keeps a local counter and periodically — every 100 ms — publishes its count and reads the fleet total. Limits are enforced against the last-known global figure.

  • Near-zero added latency.
  • Deliberately approximate: during a synchronisation interval, the fleet can overshoot by up to (servers × requests in 100 ms).
  • Right when: the limit is a soft quota and latency is critical. Wrong when the limit protects something that must not be exceeded.

The response

With counters correct, the last part of the design is what the client sees when it is limited.

What the rejected client is toldA bare 429• The client retries immediately• Every client retries together• The stampede becomes the outage429 with headers• Retry-After gives a wait time• Remaining and Reset warn early• Jittered retries spread the load
A limiter that does not tell clients how to back off manufactures the traffic spike it exists to prevent.

Return HTTP 429 Too Many Requests. Not 403 (that means "you are not allowed at all"), not 503 (that means "the server is broken"). The status code is how automated clients and libraries decide whether to retry.

Include headers so a well-behaved client can adapt:

Text
HTTP/1.1 429 Too Many RequestsRetry-After: 12RateLimit-Limit: 100RateLimit-Remaining: 0RateLimit-Reset: 12

Retry-After is a long-standing HTTP header giving seconds (or a date) until the client may try again. The RateLimit-* family has been through several draft specifications and X-RateLimit-* variants remain widespread in the wild, so pick one form and document it rather than assuming clients know it.

Send the limit and remaining headers on successful responses too. A client that can see it has 8 requests left will slow down; one that only learns at rejection cannot.

Per-tier limits

Free, paid, and internal clients should not share a limit. Store the tier's limit alongside the counter or resolve it from a small cached lookup keyed by API key. The design change is minor — the limit becomes a variable rather than a constant — but say it, because it is the difference between a limiter and a product feature.

Add a burst allowance on top of the sustained rate: token bucket gives this naturally, and it matters because real clients batch their work.

The exactness-versus-latency trade

This is the follow-up that most often decides the outcome of this question. Every request makes a synchronous call to the counter store, so:

  • Exact enforcement costs 0.5 ms per request and makes the counter store a dependency of every request. At 50,000 requests per second that is 50,000 operations per second on a sharded store, sized and monitored as a critical component.
  • Approximate enforcement (local counters synchronised every 100 ms) costs nothing per request and allows an overshoot bounded by the synchronisation interval.

The recommendation, with the condition attached: "Synchronous shared counters, because at a 200 ms latency target the 0.5 ms is 0.25% of the budget and I would rather have correct limits. If this were a 5 ms internal service, or if the counter store became the bottleneck, I would move to local counters synchronised every 100 ms and accept an overshoot of a few percent — the limit is a business policy, not an invariant."

The other follow-ups

  • Fail open or closed? Covered in The problem, and the questions to ask. Say which, and put a 10 ms timeout on the counter call so a slow store degrades to fail-open behaviour instead of stalling requests.
  • Preventing retry storms. Limited clients retry; if they all retry after exactly 12 seconds, you get a synchronised wave. Add jitter to Retry-After per client, and document that clients should back off exponentially (Reliability patterns).
  • Different costs per endpoint. A search query may cost twenty times a profile read. Charge a weighted number of tokens rather than one per request — token bucket handles this naturally and fixed-window counters do not.
  • Multi-datacentre. Global exactness would require a cross-region round trip of 100 ms+ per request, which no API can afford. The practical answer is per-region limits summing to the global limit, with an accepted overshoot when traffic is unevenly distributed. Say the overshoot out loud.
  • Monitoring. Track the rate of 429 responses per key and in aggregate. A sudden rise means either an attack or a limit that is too tight for a legitimate customer, and you cannot tell which without the per-key breakdown.