System Design Ch.5: Rate Limiters
Outline
- 0:00 One Polite Customer. One Outage.
- 0:15 Not an Attack. Just Too Many Requests.
- 1:10 Rate Limiting Is Capacity Sharing.
- 2:10 Token Bucket.
- 2:54 Sliding Window Log.
- 3:42 Sliding Window Counter.
- 4:43 One Counter. Many Servers.
- 5:46 The Hot Key Bottleneck.
- 6:55 Local Enforcement Beats Global Latency.
- 8:09 Enforce at the API Gateway.
- 8:38 User, IP, and Route Limits.
- 9:31 Rate Limits vs Concurrency Limits.
- 10:36 How to Say No.
- 11:19 Make the Limit a Contract.
- 11:46 The Rate Limiter Is a Fairness System.
- 12:04 The Architectural Blueprint.
- 12:24 Next: Distributed Cache.
Transcript
0:00 Picture the scenario right. It's three in the morning. Oh, yeah. The classic pager waking you up out of a dead sleep. Exactly. Your pager goes off. You drag yourself to your laptop, check the incident channel, and it is just flooded with stack traces. It's always a wall of red text. Yeah. And the auth service has completely exhausted its database connection pool. Legitimate users are getting timed out. The system is grinding to a halt. And your immediate assumption is, well, obviously we're under a massive DDoS attack.
0:28 Right, because that's what it looks like on the surface. But then you pull up the logs, and the truth is actually much more embarrassing. It's not an attacker at all, is it? No, not at all. It's a legitimate paying customer. They deployed a, well, a slightly buggy cron job. Oh, man. Yeah, they wrote this simple, like, five-line bash script to debug a retry bug in their own stack. And that little loop is just hammering your login endpoint as fast as their server can open network sockets. So nobody exploited a zero-day vulnerability or anything.
1:00 Exactly. The system just politely answered every single question it was asked, you know, as efficiently as it could until it physically ran out of threads. Which is such a painful way to go down. Because the overarching theme in all these postmortems we looked at is a fundamental paradigm shift. I mean, engineers frequently misclassify rate limiting as a security mechanism. Yeah, we mentally group it with web application firewalls, right? Yeah, with WAFs and DDoS mitigation. But rate limiting is actually a capacity-sharing feature.
1:27 That's the mindset shift. Your infrastructure has a strictly finite budget of database connections, worker threads, downstream network bandwidth. It's not infinite no matter what the cloud providers tell you. Exactly. So the rate limiter is the component whose sole responsibility is to say no on purpose. At a mathematically predetermined threshold, you know, to protect that finite budget. It makes sure one noisy neighbor can't starve the rest of the multi-tenant environment. Exactly. Yeah. And that protection starts with the math of counting traffic.
2:02 Because choosing the algorithm fundamentally changes the memory footprint and the fairness of your API. Let's start with the industry default, which is the token bucket. Yeah, the token bucket. Conceptually, I like to think of it as literally a prepaid transit card. You have a card that automatically tops up a few dollars every single minute, but it's capped at a maximum balance. That's a perfect analogy. Right. Like if you haven't ridden the subway all week, your card is full. You can swipe it 20 times in a row for your entire team.
2:31 You can burst your usage. But once that balance hits zero, you are stuck waiting at the turnstile for the next incremental refill. But the tradeoff there is a lack of strict precision. Yeah, you lose that exactness. The token bucket can't answer the question, exactly how many requests did this specific user execute in the last 60 seconds? It only knows the current state of the token balance. Which forces you right into the second algorithm, which is the sliding window log. Ah, the heavy one. Yeah.
2:59 To guarantee absolute precision, you have to track every single event. It's like a bouncer at a club with a clipboard writing down the exact millisecond timestamp of every single person who walks through the door. Sounds exhausting. It is. Before letting the next request through, the system has to scan the log, purge any timestamps older than the 60-second window, and then count exactly how many entries remain. The accuracy is flawless, but... But the memory footprint is catastrophic. Catastrophic. The scaling complexity is O of N, where N is your traffic volume.
3:31 The infrastructure cost of computing the limit vastly outweighs the benefit of the precision. Which forces most teams into the third option, the practical compromise. The sliding window counter. Yeah, the sweet spot. You bucket the request into fixed intervals, say, one-minute buckets, and store a simple integer counter. When a new request arrives, you look at the current minutes counter, and you add a weighted percentage of the previous minutes counter based on how much the windows overlap. Oh, I see.
4:01 So if you are 15 seconds into the current minute, that means you're 25% into the window. Exactly. So you take the current count plus 75% of the previous minutes count. Wow. So you collapse the memory cost from thousands of timestamps per user down to exactly two integers per user. Just the previous bucket and the current bucket. Yep. Two integers. You lose a microscopic fraction of precision right at the boundaries if the traffic wasn't perfectly distributed. But, well, you save an absolute fortune in memory.
4:33 It's the gold standard for production rate limiters when you need a smooth, relatively strict limit, but refuse to pay the compute overhead of a pure log. Exactly. But look, solving the algorithm on a whiteboard is entirely different from deploying it across 50 application servers behind a load balancer. Oh, yeah. Distributed state is where it gets messy. Right. Because if every server maintains its own local counters in memory, an attacker can bypass the limit by simply round-robinning their requests across all 50 nodes.
5:02 The nodes have to coordinate. They have to agree on a single source of truth for the current count. And the standard industry move here is centralized state. And that usually means Redis, right? Redis is the undisputed champion for this. You have every request, regardless of which ephemeral app server processes it, execute an atomic increment operation on a centralized Redis key. Like user colon 123 colon minute colon 42. Exactly. And because Redis is an in-memory data store that executes commands via a single-threaded event loop, an increment command is inherently atomic.
5:37 There are no race conditions. So all 50 app servers see the exact same truth without dealing with complex distributed locks. Right. But wait, the phrase single-threaded event loop introduces a massive architectural bottleneck. This is the classic hotkey problem. Yeah, it is. Think of your incoming traffic as a massive 20-lane highway. Using a single Redis key for a global rate limit is basically forcing all 20 lanes to funnel into a single toll booth. That's exactly what it is. And it does not matter if your Redis cluster has 19 other shards sitting completely idle.
6:10 Because the hash of that specific rate limit key pins it to one specific node. Yep. Which means it is bottlenecking on one specific CPU core. A single CPU core can process an incredible number of Redis operations. But, you know, it does have a ceiling. You have to split the hotkey into sibling keys. Sibling keys. Okay. So instead of incrementing user 123, the application code appends a random integer from 1 to 10, spreading the writes across user 123-1 through shard 10. Exactly. The client-side routing algorithm hashes those sibling keys across the entire Redis cluster.
6:45 So now you're utilizing all available CPU cores. You unlock the ability to utilize the entire Redis cluster for a single high -throughput limit. Okay. But taking that distributed complexity a step further brings us to the multi-region reality. Say the application is deployed in three continents for high availability. Does the Redis counter still live in just one primary region? Oh, that's a great question. And the answer is, it shouldn't. If the rate limit cluster is physically located in Virginia, a client request landing in the Tokyo Data Center has to pay a 150-millisecond cross -region round trip across the Pacific Ocean.
7:22 Just to ask permission to execute. Exactly. Imposing a 150-millisecond latency penalty on every single API request just to increment a counter totally defeats the purpose of deploying a multi-region architecture in the first place. You have to intentionally accept loose enforcement. You run an independent Redis cluster in each geographical region. The Tokyo Gateway checks the Tokyo Redis cluster. The Virginia Gateway checks the Virginia cluster. So you let the global sum drift from the true count.
7:51 Yeah, you have to. Because the goal of the rate limiter is to prevent the database from crashing. Exactly. If the database is replicated globally and can handle the aggregate capacity, localized limits are vastly superior to injecting cross-region latency into the critical path of every single request. That makes total sense. Let's pivot slightly to where this enforcement logic actually lives within the stack. Because even if we solve the math and the Redis cluster is perfectly tuned with sibling keys and regional drift, putting this logic inside the application code feels like a massive architectural anti-pattern.
8:24 Oh, it absolutely is an anti-pattern. The enforcement layer firmly belongs in the API gateway, you know, utilizing proxies like Envoy or Nginx. Right, because rate limiting is infrastructure, not a business logic library. Exactly. And the gateway is also uniquely positioned to orchestrate the overlapping dimensions of the limits. Right, because the rules are rarely singular. You can key limits based on different attributes, like authenticated traffic is keyed on the user ID, which maps cleanly to their specific pricing tier.
8:54 But unauthenticated traffic has no user ID. Which leaves you keying on the IP address. And this is incredibly dangerous because of NAT gateways. Yeah, that risk dictates why IP limits must be significantly looser and why you introduce a third dimension, which is the specific endpoint. Ah, right. Because an expensive, heavy database aggregation query requires a radically tighter budget than a lightweight health check endpoint. Exactly. The API gateway evaluates all these dimensions simultaneously.
9:24 The user, the IP, and the route. And the golden rule of enforcement is that the first limit to trip wins. Wait, if we are relying on the API gateway to drop traffic, aren't we just moving the bottleneck? Like, a massive concurrency spike will just crush the gateway's worker threads before it even reaches the backend application. Rate limiting volume over time doesn't prevent a fleet of slow connections from exhausting the gateway. That highlights the critical distinction between rate limiting and concurrency limiting.
9:54 They address entirely different failure modes. Okay, unpack that for us. So a rate limit caps volume over a time window, say 1,000 requests per minute. But a concurrency limit caps the number of requests actively in flight at a single millisecond. Oh, I see. Yeah. A user could be well under their volume budget of 1,000 per minute. But if they execute complex search queries that take 10 seconds each to return, they could hold open 200 simultaneous connections. They are respecting the rate limit, but they are hogging all the active threads.
10:23 The system still crashes. Exactly why both are required. The rate limit protects the database from sustained query volume. The concurrency limit protects the load balancer and gateway from thread exhaustion. So once the gateway determines a request violates the threshold, it has to communicate that rejection back to the client. And the mechanics of that response dictate whether the client backs off smoothly or initiates a violent retry storm. It's all about how you say no. Returning a 500 error for a rate limit violation is an operational disaster.
10:54 Because a 500 implies your system is broken. Exactly. Modern clients are programmed to handle 500s by immediately triggering automated failover scripts, which just shifts that massive wave of aggressive traffic to another availability zone, potentially cascading the outage. The system must return a 429 too many requests status code. But the status code alone isn't enough to prevent a retry storm. You have to include the headers. There are three specific headers that establish a programmatic contract here.
11:24 Yep. X-RateLimit-Remaining exposes the exact number of requests left in the current window. Then X-RateLimit-Reset provides the exact Unix timestamp when the window rolls over. And finally, on the actual 429 rejection payload, you include the retry after header, explicitly instructing the client how many seconds they must wait before establishing a new connection. Let's summarize the mental model we've constructed today. The algorithm dictates the tradeoff between precision and compute cost. The token bucket models bursty realism.
11:55 The sliding window log provides heavy audited precision. And the sliding window counter delivers the cheap, practical middle ground. Exactly. Enforcement strictly belongs in the API gateway to drop requests before they consume backend threads. Redis handles the distributed counting state, but you have to architect around the single-threaded hotkey bottleneck using sibling keys. And you have to tolerate regional drift. Yep. Embrace the drift. Finally, using a 429 status code paired with explicit headers turns your clients into cooperative partners rather than blind attackers.
12:28 That's it for today's chapter. See you in Chapter 6.