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
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.
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:
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:
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.
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:
HTTP/1.1 429 Too Many RequestsRetry-After: 12RateLimit-Limit: 100RateLimit-Remaining: 0RateLimit-Reset: 12Retry-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-Afterper 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.