skip to content

A team moves rate limiting from each backend service into the shared API gateway in front of them. What does the gateway need to track to enforce limits correctly, and what breaks if the gateway is horizontally scaled to multiple instances without care?

level: seniorimportance: must knowfreq 65%

answer

  1. counter per client per window
  2. token bucket vs sliding window
  3. local memory undercounts across N instances
  4. shared Redis store for correctness
  5. fail open vs fail closed on store outage

basics

~20 s

The gateway counts how many requests each client has made in a time window and blocks them once they go over the limit. If there are several gateway machines each counting on their own, a client can sneak past the real limit by spreading requests across them.

solid answer

~50 s

Centralizing rate limiting at the gateway means one component tracks request counts per client (by API key, user ID, or IP) over a rolling or fixed time window and rejects requests once a threshold is exceeded, typically with a 429 status and a Retry-After header. Common algorithms are token bucket (allows bursts up to a bucket size, refills over time) and sliding window (smooths counting across window boundaries). The critical scaling problem is that a single gateway instance can track counts in local memory cheaply, but once you run multiple gateway instances behind a load balancer, each instance only sees a fraction of a client's traffic, so local counters undercount and the real per-client rate exceeds the configured limit. The standard fix is a shared, low-latency store, usually Redis, that all gateway instances read and write counters against atomically, trading a small amount of added latency and a new dependency for correctness. Some setups accept approximate limiting (each instance enforces limit/N) to avoid that shared-state cost.

go deeper

for a junior

Should know rate limiting caps requests per client over time and rejects excess requests with an error.

for a middle

Should be able to name token bucket or sliding window and explain roughly how each avoids simple counting flaws.

for a senior

Should identify the multi-instance undercounting problem unprompted and describe the shared-store fix along with its cost (latency, new dependency, race conditions needing atomic ops).

for a principal

Should reason about the fail-open/fail-closed trade-off explicitly, weigh shared-store precision against per-instance approximate limiting at very large scale, and connect this to broader capacity-protection strategy across the platform, not just the gateway.

## Why it moves to the gateway **Rate limiting** is the practice of capping how many requests a given client can make in a period of time, protecting backend capacity from being overwhelmed by a single noisy client, whether that's a buggy retry loop, a scraper, or a deliberate denial-of-service attempt. Offloading it to the gateway, rather than implementing it inside every backend service, is attractive for the same reason authentication and TLS offloading are: it's a cross-cutting concern that behaves identically regardless of which service is being called, so implementing it once avoids inconsistent per-service limits and duplicated logic. ## What the gateway has to track Mechanically, the gateway identifies the caller, usually by API key, authenticated user ID, or source IP as a fallback, and maintains a **counter** tied to that identity plus a time window. Two algorithms dominate in practice. | Algorithm | How it counts | |---|---| | **Token bucket** | assigns each client a bucket that holds up to N tokens; each request consumes one token, and tokens refill at a steady rate over time, which naturally allows short bursts up to the bucket size while enforcing a steady average rate over the long run | | **Sliding window** | counts requests in overlapping time slices to avoid the classic fixed-window problem where a client can send double the intended limit by timing requests to straddle a window boundary (for example, bursting at the end of one minute and the start of the next) | Whichever algorithm is used, once a client exceeds its allotment, the gateway rejects the request immediately, typically with an HTTP `429 Too Many Requests` status and a `Retry-After` header telling the client how long to back off, without the request ever reaching, or costing capacity from, any backend service. ## The scaling problem The scaling problem emerges the moment the gateway itself is no longer a single process. A single gateway instance can track counters cheaply in local memory: fast, no network hop, no extra infrastructure. But production gateways are almost always horizontally scaled behind a load balancer for availability and throughput, meaning a given client's requests get distributed across several gateway instances, often somewhat randomly depending on the load balancing algorithm. If each instance keeps its own local counter, none of them individually see the client's full request volume, so a client sending 1000 requests split evenly across 5 gateway instances looks like only 200 requests to each instance, and a limit set at 500/instance would let 1000 through when the intended limit was 500 total. This **under-counting** is the central failure mode: rate limiting appears to be working (no errors, counters look fine) while the actual aggregate rate silently blows past the configured ceiling, often only discovered when a backend gets overwhelmed anyway despite "working" rate limiting. ## The standard fix and what it costs The standard fix is to move the counter state out of each gateway instance's local memory and into a shared, low-latency store that every instance reads and writes atomically, almost always **Redis**, using atomic increment-and-check operations (or Redis's built-in Lua scripting) to avoid race conditions where two instances both read a stale count and both allow a request that together exceeds the limit. This restores correctness but introduces a new dependency and a small amount of added latency per request, plus a new failure mode of its own: if the shared store becomes slow or unavailable, the gateway has to decide whether to - **fail open** (allow all requests, risking overload), or - **fail closed** (reject all requests, risking a full outage over what should be a soft-limiting feature), and that choice needs to be deliberate, not accidental. Some teams accept an approximate compromise instead: divide the intended global limit by the number of gateway instances and enforce that smaller limit locally per instance, trading precision (the real limit can drift somewhat above or below intent as instance count or traffic distribution changes) for avoiding the shared-store dependency entirely. ## Where it shows up in practice A well-known real-world instance of this pattern is Kong or Envoy's rate-limiting filter backed by a Redis cluster, and cloud API gateways like AWS API Gateway or Azure API Management, which expose managed, distributed rate limiting specifically because customers running these gateways at scale hit exactly this local-counter-undercounts problem when they first tried to DIY it with in-memory counters on auto-scaled gateway fleets.

  • Why does a fixed time window (e.g. reset every 60 seconds) allow a client to burst past the intended rate?
    A client can send the full limit right at the end of one window and the full limit again right at the start of the next, achieving roughly double the intended rate within a short span that straddles the boundary, while each individual window's count still looks compliant. A sliding window or token bucket avoids this by smoothing the count across overlapping time or continuous refill instead of resetting sharply at fixed boundaries.
  • If the shared rate-limit store (like Redis) becomes unavailable, should the gateway fail open or fail closed, and why does it matter?
    Failing open (allowing all requests through) risks the exact overload scenario rate limiting was meant to prevent, while failing closed (rejecting all requests) turns a soft-limiting feature outage into a full service outage for every client, including well-behaved ones. Most production systems choose to fail open for short store blips with monitoring/alerting, since a brief unlimited window is usually less damaging than a full outage, but this is a deliberate trade-off that should be documented, not a default left to chance.
  • Why might a team accept per-instance approximate limits (dividing the global limit by instance count) instead of a shared store?
    It avoids adding a new stateful dependency, extra network hop latency, and a new failure mode (the shared store itself becoming a bottleneck or single point of failure) to every rate-limited request. The trade-off is that the effective limit becomes imprecise, drifting as instance count or traffic distribution across instances changes, which is acceptable when the limit is a rough protective ceiling rather than a hard contractual guarantee.

Like several separate ticket booths at a stadium each independently deciding whether a fan has bought too many tickets today, without comparing notes with the other booths; a scalper just walks between booths and gets past every individual booth's limit even though the total across all booths is way over the real cap.

saying these in an interview costs you the question

  • assumes local in-memory counters work correctly once the gateway is scaled horizontally
  • doesn't know why fixed windows allow boundary bursting
  • can't name a mechanism (shared store, approximate division) to fix cross-instance undercounting
  • has no opinion on fail-open vs fail-closed behavior when the counter store is down
  • conflates rate limiting with authentication or authorization

context