skip to content

You're designing a rate limiter for an API gateway deployed across three geographic regions, each with many stateless gateway instances, enforcing one global per-API-key limit. A single centralized Redis in one region adds cross-region latency to every request and becomes a single point of failure. Walk through the design trade-offs: what happens if you use per-region local counters instead of a shared store, and what happens if the central Redis becomes unreachable?

level: principalimportance: should knowfreq 45%

answer

  1. local counters miss cross-region usage
  2. static partition vs periodic sync
  3. GCRA = token bucket without discrete windows
  4. circuit breaker + fail-open fallback
  5. accuracy vs latency is the real trade-off

basics

~30 s

Local per-region counters are fast but only see local traffic, so a client hitting all three regions can get up to 3x its real limit unless you divide the budget per region. A central Redis is accurate but adds latency and is a single point of failure — you must decide in advance whether to fail open (allow everything, risk overload) or fail closed (reject everything, cause an outage) when it's unreachable.

solid answer

~60 s

Per-region local counters avoid a network hop to a remote store on every request, which matters at gateway scale, but each region only sees the traffic that landed there — a client rate limited to 100/sec globally but invisible across regions can send 100/sec to each of three regions and get 300/sec effectively. Common mitigations are statically dividing the global budget across regions (each region enforces 33/sec), or periodically syncing regional counts asynchronously to approximate the global total with some lag and slop. A centralized store gives an exact global count but adds a network round trip's worth of latency to every request and turns Redis into a single point of failure: you must explicitly choose whether an unreachable Redis fails open (allow all traffic — risks the very overload the limiter exists to prevent) or fails closed (reject all traffic — turns a rate-limiter outage into a total API outage), and most production systems fail open with a tight circuit breaker plus a fallback to local approximate limiting rather than either extreme.

go deeper

for a junior

Not generally expected to reason about multi-region trade-offs; can note in general terms that a single shared server 'far away' would be slow and could go down.

for a middle

Should recognize that per-region counters miss traffic from other regions and that a centralized store adds latency, even without proposing a specific fix.

for a senior

Expected to propose at least one concrete mitigation (static partitioning or periodic sync) and reason about the accuracy/latency trade-off explicitly.

for a principal

Expected to design the full picture: partitioning or sync strategy, an explicit fail-open/fail-closed decision backed by a circuit breaker and fallback, and to articulate that this is a spectrum, not a single 'correct' answer, informed by what's been operated at scale.

## Two compounding problems At gateway scale — many stateless instances across multiple regions, high request-per-second volume — the naive "every request checks a single centralized Redis" design runs into **two compounding problems**, and understanding both is what separates a limiter that survives production from one that becomes the next incident. ## Latency, and the gap local counters open The first problem is **latency**: if Redis lives in one region, every request originating in a different region pays a cross-region round trip, commonly tens of milliseconds, purely for the rate-limit check, before any actual business logic runs. At high request volumes this either becomes a meaningful tax on p99 latency for every single request, or forces batching/pipelining tricks that add their own complexity and staleness. The mechanism-level fix people reach for first is **per-region local state**: each region's gateway instances check a rate limiter backed by local memory or a regional Redis, incurring no cross-region hop. This is fast, but it creates a correctness gap: - the whole point of a global per-API-key limit is that it should reflect the client's total usage across all regions, - a purely local counter has zero visibility into what the same client is doing in the other two regions. A client that discovers, deliberately or by accident (e.g. a multi-region deployment of their own client fleet), that requests get load-balanced across regions can send up to N times their intended limit, where N is the number of independent regional counters, simply by spreading traffic across them. ## The mitigations, and what each one buys The standard mitigation is to accept an explicit accuracy/latency trade-off rather than pretend both can be had perfectly. 1. **Static budget partitioning** — divide the global limit across regions in proportion to expected traffic share, e.g. a 300/sec global limit becomes 100/sec enforced independently in each of three regions, which is simple and requires zero cross-region communication, but is wasteful when traffic is uneven: a region with little traffic wastes its allotted budget while a hot region's clients get throttled even though global capacity remains. 2. **A more adaptive approach** periodically, e.g. every 1-5 seconds, aggregates regional counts to a shared view and rebalances allotments, trading some staleness — a client can transiently exceed the true global limit within that sync interval — for much better utilization than static partitioning, at the cost of running an additional aggregation/sync pipeline that is itself infrastructure to build, monitor, and fail-handle. 3. **A third pattern** avoids the discrete-window problem altogether by using GCRA (Generic Cell Rate Algorithm, the algorithm behind the popular throttled and redis-cell libraries) — mathematically equivalent to token bucket but tracked as a single "theoretical arrival time" value per key rather than a counter, which composes cleanly with a centralized store and gives smooth, precise rate control without discrete window boundary artifacts, but it still requires the same centralization-vs-locality decision as any other shared-state approach. ## When the shared store is unreachable The second, orthogonal problem is what happens when the centralized store itself becomes unreachable — a Redis failover, a network partition between a region and Redis's home region, or Redis simply falling over under load. This decision must be made explicitly and in advance, because both extremes are dangerous. - **Failing closed** — rejecting all requests when the limiter can't be consulted — converts a rate-limiter dependency outage into a full API outage, which is almost always the wrong trade for a component whose entire purpose is protective, not core functionality; you've let a supporting system become as critical as the primary path. - **Failing open** — allowing all requests through unchecked — removes exactly the protection the limiter exists to provide at precisely the moment (Redis under stress, likely correlated with a broader incident) when uncontrolled traffic is most dangerous. The production-grade answer used by most large-scale gateways is neither pure extreme: wrap the Redis dependency in a **circuit breaker** with a short timeout, and on trip, fall back to degraded local-only limiting, each instance enforcing a conservative local approximation of its share of the global budget, for the duration of the outage, then resume centralized checks once Redis recovers — accepting slightly worse fairness/accuracy during the incident window in exchange for neither a full outage nor unlimited exposure. ## The deeper principle The deeper principle this scenario illustrates is that at sufficient scale, exact global rate limiting and low added latency are in direct tension, and every real design — static partitioning, periodic sync, GCRA with regional replicas, circuit-broken fallback — is a different point on that trade-off curve rather than a way to escape it, which is exactly why this question separates candidates who've only implemented a single-instance limiter from those who've operated one at scale.

  • How does GCRA avoid the discrete window boundary problems that fixed and sliding window counters have?
    GCRA tracks a single continuous value per key — the theoretical arrival time (TAT) of the next allowed request — rather than a bucket that fills and resets on a timer. Each request compares 'now' against TAT and either advances TAT (if allowed) or rejects, so there's no clock-aligned reset boundary to exploit the way there is with fixed-window counters.
  • What's a concrete symptom that would tell you a multi-region limiter is using purely local per-region counters with no cross-region visibility?
    A client that spreads its own traffic evenly across the three regions, deliberately or via normal geo-load-balancing, is able to sustain a rate roughly N times its documented single-region limit, where N is the number of regions — the giveaway is that the effective global rate scales with the number of active regions rather than staying fixed.
  • Why might failing open with a degraded local fallback still be risky even though it avoids a full outage?
    During the fallback window, the limiter is only enforcing its local, non-global view, so a client that discovers the outage, or is simply distributed across regions, can exceed its true global limit for the duration of the Redis outage — an acceptable, bounded risk in most designs, but one that should be monitored and alerted on rather than silently tolerated indefinitely.

It's like three separate bank branches each keeping their own ledger of a customer's withdrawals instead of one shared account balance — the customer can withdraw their limit at each branch separately unless the branches sync up, and if the head office (central ledger) goes down, each branch must decide on the spot whether to keep honoring withdrawals from memory or freeze all of them.

saying these in an interview costs you the question

  • Assumes a single centralized Redis has no downsides beyond 'it's a single point of failure'
  • Doesn't recognize that local-only counters undercount or miss cross-region traffic
  • Picks fail-open or fail-closed without acknowledging the trade-off the other way
  • Thinks periodic sync gives the same accuracy as a fully centralized check with zero added latency

context