Marcio Cunha

Distributed Rate Limiting in Microservices with Token Bucket and Redis Cluster

Learn how to design distributed traffic control in microservice systems using the token bucket algorithm and coordinated Redis Cluster instances to mitigate overloads and denial-of-service attacks.

Marcio Cunha•5 min
Also available in:PortuguêsEspañol
Summary
  • The token bucket algorithm allows controlled traffic spikes while ensuring a steady average of requests per second.
  • The Redis Cluster architecture avoids single-point-of-failure bottlenecks by partitioning control data across multiple nodes.
  • Lua scripts executed atomically in Redis eliminate race conditions during token verification and consumption.
  • Local fallback strategies ensure the system remains operational even when the caching cluster experiences downtime.
  • Keys composed of user identifiers and time windows prevent incorrect global exhaustion of shared resources.

The Challenge of Protecting Microservices Against Excessive Traffic

In modern microservices architectures, elastic scalability is both a blessing and a curse. While we easily add new server instances to handle traffic spikes, our relational databases, message queues, and third-party APIs continue to have strict capacity limits. This is precisely where rate limiting comes into play: a fundamental software engineering technique that acts like a strict bouncer at the door of a busy venue, controlling exactly how many requests each client can make within a given time interval to prevent cascading failures.

When running a monolithic application on a single server, tracking accesses is a trivial task maintained in local RAM. However, when we distribute our application across dozens of Docker containers orchestrated by Kubernetes, the scenario changes dramatically. If a malicious user or a client with a software bug sends thousands of requests per second, these calls might be distributed randomly across different nodes of our infrastructure. Without a centralized and coordinated counting mechanism, each server will mistakenly assume traffic is low, allowing the total request quota to be exceeded multiple times.

Understanding the Token Bucket Algorithm in Practice

There are several mathematical ways to limit traffic, but the token bucket stands out as the industry gold standard due to its unmatched flexibility. Imagine a bucket that stores tokens up to a predefined maximum limit. An invisible faucet pours new tokens into this bucket at a constant rate, for example, ten tokens per second. Each time a client makes a request to our API, the system attempts to remove one token from that client's corresponding bucket. If the bucket has enough tokens, the request is authorized immediately and the token is discarded. If the bucket is completely empty, the request is rejected with the famous HTTP status code 429 Too Many Requests.

In practice, this means the token bucket solves one of the biggest problems of traditional counters: it tolerates legitimate traffic bursts. If a mobile app needs to load ten images simultaneously upon opening, the full bucket allows all of them to pass at once, provided the long-term average respects the replenishment rate. Conversely, stricter algorithms would reject access immediately after the first exceeding request of that microsecond. To implement this efficiently in distributed systems, we need high-speed storage that supports extreme concurrency without losing temporal precision.

Redis Cluster Architecture for High Availability

To synchronize the state of token buckets across dozens of application servers, we need an ultra-fast in-memory database, and Redis emerges as the community's natural choice. More than a simple isolated instance, Redis Cluster offers automatic data sharding and native replication, dividing keys across up to 16,384 hash slots distributed among multiple nodes. This means that even if one cluster node fails during a traffic surge, the rest of the infrastructure remains operational, ensuring the resilience required in mission-critical environments.

However, operating distributed counters in Redis introduces subtle concurrency pitfalls. If two application servers read the remaining token count in the same millisecond, both calculate that there is still space and decrement the value separately, resulting in a race condition that corrupts the limit's accuracy. To shield our system against this type of failure, we use Lua scripts executed directly on the Redis server. Since Redis processes commands and scripts completely atomically, we guarantee that reading, verifying, and updating the bucket happen as a single, indivisible transaction.

Practical Implementation of the Algorithm with Lua Scripts

Below we present a functional example of a Lua script designed to run in Redis, implementing the logic of gradual replenishment and token consumption per client key. This script calculates the elapsed time since the last request, adds newly generated tokens based on the configured rate, and validates whether there is enough balance to authorize the current transaction.

local key = KEYS[1]local now = tonumber(ARGV[1])local capacity = tonumber(ARGV[2])local fill_rate = tonumber(ARGV[3])local requested = tonumber(ARGV[4])local data = redis.call('HMGET', key, 'tokens', 'last_updated')local tokens = tonumber(data[1])local last_updated = tonumber(data[2])if not tokens then    tokens = capacity    last_updated = nowelse    local delta = math.max(0, now - last_updated)    tokens = math.min(capacity, tokens + delta * fill_rate)endlocal allowed = 0if tokens >= requested then    tokens = tokens - requested    allowed = 1endredis.call('HMSET', key, 'tokens', tokens, 'last_updated', now)redis.call('EXPIRE', key, math.ceil(capacity / fill_rate))return {allowed, tokens}

To integrate this script into the backend application, we send dynamic parameters with every HTTP request received at the API gateway. If the return indicates the operation was permitted, the flow proceeds normally to the internal microservices. Otherwise, we interrupt execution immediately, saving valuable computing resources on backend servers and protecting the ecosystem against destructive overloads.

Fallback Strategies and Production Fault Handling

No distributed system is immune to network outages or temporary infrastructure failures, and blindly relying on Redis Cluster can turn into a single point of catastrophic failure if the cache becomes unreachable. When the Redis cluster suffers an outage, the worst possible scenario is causing all client requests to fail due to the inability to validate rate limiting. To mitigate this operational risk, mature architectures implement intelligent fail-open or local in-memory fallback mechanisms.

In practice, this means that if the Redis call returns a timeout or connection refused error after a strict millisecond limit, the limiting middleware assumes temporary permissive behavior or falls back to a local in-memory heuristic counter. This decision ensures business continuity for legitimate users during infrastructure incidents, while automated alerts notify the engineering team to restore the health of the cache cluster as quickly as possible.

Final Considerations on Scalability and Resilience

Implementing a robust distributed rate limiting system using the token bucket algorithm and Redis Cluster requires a careful balance between mathematical precision, network latency, and fault tolerance. Although the infrastructure cost to maintain dedicated cache nodes exists, it is infinitely smaller than the financial and reputational damage caused by the unavailability of entire platforms due to a lack of protection against traffic spikes. By adopting atomic Lua scripts and defensive fallback strategies, engineers can build highly resilient microservice ecosystems capable of absorbing request storms without losing composure.