Skip to main content

Command Palette

Search for a command to run...

Distributed Rate Limiting

The Real Challenge Isn't Counting Requests

Updated
•7 min read•View as Markdown
Distributed Rate Limiting

Most discussions around rate limiting start with algorithms — Fixed Window, Sliding Window, Token Bucket, Leaky Bucket.

After building and operating a distributed rate limiting layer, one thing becomes clear quickly:

The algorithm is the easy part. The hard part is knowing how many requests a client has already made — consistently — across every node that might serve their next request.


Why rate limiting exists

Every service exposed to the outside world eventually encounters traffic it didn't plan for. Normal users. Automated integrations. Retry storms from buggy clients. Batch jobs that hit at midnight. Occasionally, something worse.

Without a control boundary, that traffic competes with legitimate requests for the same threads, connections, and downstream capacity.

Rate limiting isn't about rejection. It's about keeping the service available for the requests that should succeed.

The shape of the problem looks simple at one instance. It stops looking simple the moment you deploy more than one.


The distributed systems problem

A single instance can keep its counter in memory. It sees every request. It has a complete picture.

Add a load balancer and two replicas, and that complete picture fractures.

Instance 1 has seen 47 requests.

Instance 2 has seen 31.

The client has actually made 78. Both instances, looking only at local state, would allow the next request. Both would be wrong.

The question the rate limiter must answer — "how many requests has this client already made?" — cannot be answered correctly by any single replica in isolation. It requires a shared view of state that doesn't exist in local memory.

That shared view is the actual design problem.


Three dimensions of limiting

Rate limiting on request count alone is a blunt instrument. In practice, the limits that matter are more granular — and they operate simultaneously on different dimensions. There can be multiple dimensions to it.

IP-based limiting protects against unauthenticated traffic, scrapers, and clients that don't identify themselves. A single IP hammering the API gets throttled regardless of what it's doing.

Tenant/API key limiting enforces the contractual boundary. A tenant on a standard tier has a different quota than one on a premium tier. This limit is about fairness across the fleet — one busy tenant shouldn't degrade everyone else's experience.

Endpoint-based limiting protects expensive operations independently of overall request volume. A search endpoint that runs a full-text query might have a limit of 20 requests per minute even for premium tenants, because the downstream cost per request is orders of magnitude higher than a lightweight read.

A request must pass all three checks to proceed. Passing two out of three isn't enough.


Why Redis fits this problem

Each dimension needs an atomic, distributed counter with automatic expiry. Redis provides exactly that.

INCR is atomic. Multiple application instances can safely update the same counter simultaneously without explicit locking — Redis serialises the operations internally. The first INCR on a key creates it with value 1; a TTL is set immediately after to define the window. When the window expires, Redis removes the key automatically. No cleanup scheduler, no background jobs.

Three counter keys per request, each independently managed:

ip:203.0.113.42        → value: 14  TTL: 47s
tenant:acme-corp       → value: 847 TTL: 47s
endpoint:search        → value: 8   TTL: 47s

The application doesn't remember previous requests, active clients, or cleanup schedules. Every request performs the same sequence: increment, evaluate, decide. The storage layer manages everything else.


The fixed window trade-off

The implementation above uses a fixed window — the counter resets at the start of each interval. Simple, cheap, and effective for most workloads.

The edge case worth knowing: at window boundaries, a client can effectively double their allowed rate.

For most use cases — protecting a backend from runaway clients, enforcing tier limits — this boundary burst is acceptable. The client still can't sustain it; the next second of window 2 they're already at 100.

For cases where the boundary burst is unacceptable — financial APIs, high-cost operations — a sliding window eliminates it by evaluating the last N seconds at the moment of each request rather than since the last reset. The trade-off is higher Redis complexity: a sorted set per client instead of a simple counter, with a range delete on each check. More expensive, more correct.


What happens when Redis is unavailable

A rate limiter that takes down your service when its backing store goes offline has inverted its purpose.

The two policies and their implications:

Policy Behavior When to use
Fail open Allow all requests when Redis is unavailable When availability outweighs risk — most internal services
Fail closed Reject all requests when Redis is unavailable When the downstream is expensive enough that uncontrolled access is worse than an outage

Fail open is the right default for most services. A brief Redis outage should not trigger a full service outage. Log the failure, emit a metric, alert — but let traffic through.

What fail open does require:

The Redis client must have aggressive timeouts configured.

A connection attempt that hangs for 5 seconds before failing means every request waits 5 seconds before being allowed through.

The rate limiter becomes a latency spike injected into every request during an outage. Set connection timeout and command timeout in the low hundreds of milliseconds.


The response tells the client what to do next

A 429 with no context forces the client to guess. Most will retry immediately, which makes the problem worse.

The standard headers give clients the information they need to back off correctly:

HTTP/1.1 429 Too Many Requests

X-RateLimit-Limit: 1000
X-RateLimit-Remaining: 0
X-RateLimit-Reset: 1719475200

Retry-After: 47

Retry-After is the most important. It tells the client exactly how many seconds to wait before trying again. Well-behaved clients respect it. The retry storm that would have followed a bare 429 becomes an orderly queue of clients reconnecting when their window resets.

When returning 429, also include which dimension was exceeded. An IP-limited response and a tenant-quota-exceeded response both return 429 — but a tenant who knows they've hit their tier limit can contact support. A tenant who gets a cryptic 429 opens a bug report instead.


The algorithm determines how. The architecture determines whether.

An elegant sliding window algorithm running per-instance with local state gives you accurate limits on one replica and wrong limits everywhere else.

A simple fixed window backed by Redis gives you slightly approximate limits that remain correct across the entire fleet.

At scale, approximate and correct beats precise and broken.

The real decisions in distributed rate limiting are architectural:

where does state live, how do you handle backing store failures, what do you tell clients when you reject them, and how do multiple limit dimensions compose into a single allow/deny decision.

The algorithm sits on top of those answers, not underneath them.


Connect on LinkedIn · Read more on Hashnode