A rate limiter is a small idea — reject requests past a threshold — that turns into a genuinely tricky distributed systems problem the moment it has to run correctly across more than one server. This is the design for a rate limiter sitting in front of a large API gateway, protecting both the gateway itself and the backends behind it.
Problem
We’re building a rate limiter that sits on the request path of an API gateway serving roughly 500,000 requests/second across 100,000 distinct API clients (identified by API key or user id). Its job is to enforce per-client limits (e.g., “1,000 requests/minute” on a free tier, “50,000 requests/minute” on enterprise) so that one noisy or abusive client can’t degrade service for everyone else, and so backends that were sized for average load don’t fall over under an unexpected burst.
The hard part isn’t the algorithm — counting requests against a threshold is conceptually trivial. It’s making that counting accurate and fast when the counting has to happen consistently across dozens of stateless gateway instances handling the same client’s traffic in parallel, without the rate limiter itself becoming the bottleneck it’s meant to prevent.
Requirements
Functional
- Enforce per-client request limits over a configurable time window (e.g., 1,000/min, 10,000/hour).
- Support multiple concurrently active limit tiers per client (a short burst limit and a longer sustained limit).
- Return a clear rejection response (HTTP 429) with retry guidance when a limit is exceeded.
- Support limit configuration changes (a client upgrades tier) taking effect promptly, without a deploy.
- Allow different limiting strategies for different endpoints (a cheap read endpoint and an expensive write endpoint shouldn’t share one counter).
Non-functional
- Added latency per request under 5ms at p99 — the rate limiter sits on every single request, so its own overhead has to be close to invisible.
- Must work correctly across a fleet of stateless gateway instances — a client hitting different instances on different requests must still be limited against one shared count, not per-instance counts that each allow the full limit.
- Available even if the limiting store has a brief outage — fail in a way that doesn’t take the whole gateway down.
- Memory-efficient at 100,000 tracked clients with sub-minute windows — this is a lot of small counters updated very frequently.
Capacity Estimation
Assume 500,000 requests/second gateway-wide, across 100,000 unique clients, with a sliding-window rate check on every single request.
- Rate limiter check QPS: matches gateway QPS exactly, since every request needs a check — 500,000/s. This is the number that rules out anything involving a network round trip to a database per request; even Redis, sub-millisecond as it is, becomes a real cost at this volume if not batched or colocated carefully.
- Counter storage: with a sliding-window-counter approach (detailed below), each client needs roughly two integers per active limit tier (current window count, previous window count) — at 2 tiers per client (burst + sustained) × 100,000 clients × 2 counters × 8 bytes ≈ 3.2 MB. This comfortably fits in memory on a single Redis node, let alone a cluster — the working set here is tiny relative to typical cache sizing.
- Redis ops/sec if centralizing all checks there: 500,000/s of
INCR-plus-EXPIREstyle operations is within reach of a well-provisioned Redis cluster (single-digit millions of ops/sec across a sharded cluster), but it means every gateway request now has a hard dependency on a network hop to Redis in its critical path — this is the number that motivates the local-cache-plus-async-sync design in Deep Dive. - Local in-memory check cost (if we avoid the network hop for the common case): a hash map lookup and increment, sub-microsecond — three to four orders of magnitude cheaper than a network round trip, which is why the architecture leans hard on doing as much locally as correctness allows.
- Rate checks / sec
- 500,000/s
- one check per gateway request
- Tracked clients
- 100,000
- 2 limit tiers each
- Counter memory
- ~3.2 MB
- 2 tiers x 100K clients x 2 counters x 8B
- Central store ops (naive)
- 500,000/s
- 1 network hop per request if not localized
- Local check cost
- sub-microsecond
- vs ~1ms+ network round trip
API Design
The rate limiter is a library/sidecar invoked internally by the gateway on every request, plus a small admin API for tier configuration. It’s not typically a public-facing API by itself, but it’s useful to model it as one for clarity.
POST /internal/ratelimit/check HTTP/1.1
Content-Type: application/json
{
"clientId": "key_8f2c19",
"endpoint": "POST /v1/orders",
"cost": 1
}
HTTP/1.1 200 OK
Content-Type: application/json
{
"allowed": true,
"limit": 1000,
"remaining": 812,
"resetAt": "2026-08-25T14:31:00Z"
}
POST /internal/ratelimit/check HTTP/1.1
Content-Type: application/json
{
"clientId": "key_8f2c19",
"endpoint": "POST /v1/orders",
"cost": 1
}
HTTP/1.1 429 Too Many Requests
Retry-After: 8
Content-Type: application/json
{
"allowed": false,
"limit": 1000,
"remaining": 0,
"resetAt": "2026-08-25T14:31:00Z"
}
PUT /internal/ratelimit/tiers/key_8f2c19 HTTP/1.1
Content-Type: application/json
{
"burst": { "limit": 100, "windowSeconds": 10 },
"sustained": { "limit": 1000, "windowSeconds": 60 }
}
HTTP/1.1 200 OK
The cost field lets expensive endpoints (a bulk export, a search query) consume more than one unit per call against the same budget as cheap endpoints, without needing entirely separate limit tracking — this is a small addition that makes the limiter meaningfully more useful in practice.
Data Model
Configuration is small and relational; counters are high-churn and belong in a fast key-value store, not a table.
CREATE TABLE rate_limit_tiers (
client_id VARCHAR(64) PRIMARY KEY,
tier_name VARCHAR(32) NOT NULL,
burst_limit INT NOT NULL,
burst_window_s INT NOT NULL DEFAULT 10,
sustained_limit INT NOT NULL,
sustained_window_s INT NOT NULL DEFAULT 60,
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
Counter state, per client per tier, lives in Redis as a small set of keys rather than SQL rows:
| Key pattern | Type | Purpose |
|---|---|---|
rl:{clientId}:{tier}:cur |
Integer, TTL’d | Count in the current window |
rl:{clientId}:{tier}:prev |
Integer | Count in the previous window, used for the sliding-window weighted estimate |
rl:{clientId}:{tier}:window_start |
Timestamp | Anchors which window cur belongs to |
rate_limit_tiers is read-heavy and changes rarely (a tier upgrade is a human or billing-system action, not a per-request event), so it’s cached in each gateway instance’s local memory with a short TTL and invalidated via a pub/sub message on update — this keeps the 500,000/s hot path from ever touching the configuration store directly.
High-Level Architecture
Every gateway instance keeps a local, in-memory approximate counter per client, checked on the hot path with no network call. Those local counters periodically reconcile against a shared Redis-backed source of truth, so no single instance can silently allow a client far past their limit just because it’s only seeing a fraction of that client’s total traffic.
Deep Dive
Algorithm choice
Four standard algorithms, and the trade-offs between them are the crux of this design:
Fixed window counter: increment a counter keyed by (clientId, currentMinute), reject once it exceeds the limit, reset at the minute boundary. Simple and cheap, but has a real correctness problem at the boundary — a client can send the full limit in the last second of one window and the full limit again in the first second of the next, getting 2x their stated limit in a 2-second span.
Sliding window log: store a timestamp for every request in a sorted set, count entries within the trailing window on each check. Perfectly accurate, but storage grows with request volume rather than staying constant — at 500,000/s this is untenable without aggressive trimming, and even trimmed, it’s a much heavier per-request operation than a counter increment.
Sliding window counter: keep two fixed-window counters (current and previous) and compute a weighted estimate: previous_count × (1 - elapsed_fraction_of_current_window) + current_count. This smooths the boundary problem from the fixed-window approach down to a bounded approximation error, at the storage cost of a plain counter (two integers, not a growing log).
Token bucket: each client has a bucket that refills at a steady rate up to a cap; each request consumes a token, and a request with no tokens available is rejected. This naturally supports bursting (a client can spend a large accumulated bucket in one burst) while still enforcing an average rate over time, which fixed and sliding window counters don’t do as elegantly.
Decision
Sliding window counter for sustained limits, token bucket for burst limits
Two algorithms, one per tier, rather than a single unified approach
Sustained limits (1,000/min) care about smoothing out the fixed-window boundary problem cheaply — the sliding window counter’s bounded approximation error is an easy trade for its low, constant memory cost. Burst limits (100/10s) exist specifically to allow short bursts while still capping average rate, which is exactly what a token bucket is built for. Using one algorithm for both would either make sustained limiting needlessly bursty or make burst limiting needlessly rigid — the two tiers are solving different problems and deserve different tools.
Local caching to avoid a network hop per request
At 500,000 checks/second, a network round trip to Redis on every request — even at roughly 1ms — adds real, avoidable latency and caps throughput at whatever Redis’s connection and ops capacity allows, rather than the gateway’s own capacity. The design instead keeps a local, per-instance approximate counter for each active client, incremented synchronously in-process (a hash map, effectively free), and only periodically (every 1-2 seconds) reconciles that local count against the shared Redis counter via a background sync worker.
This means any single gateway instance decides against a count that’s slightly stale relative to the true global count — a client spread evenly across 20 instances could theoretically get up to roughly 20x their limit through in a worst case if every instance’s local view lags simultaneously. In practice this is bounded and rare: the sync interval is short, the local counter enforces limit / instanceCount as a soft per-instance cap between syncs, and the product accepts a small amount of soft overshoot as the cost of not putting Redis on the hot path.
Where the rate limiter sits in the request path
It runs as a gateway-level middleware, before the request is routed to any backend — rejecting an over-limit request as early and cheaply as possible is the entire point, since the alternative (letting it reach a backend service and rejecting there) wastes the exact capacity the limiter exists to protect.
Per-endpoint limits and cost weighting
Not every endpoint costs the backend the same amount of work. A GET /v1/status call and a POST /v1/reports/generate call shouldn’t share a naive per-request counter, because a client could exhaust their “request budget” entirely on cheap calls and never be constrained on expensive ones, or the reverse — a client doing legitimate bulk work gets blocked by a limit sized for lightweight traffic. The cost field in the check request lets each endpoint declare its weight against the shared budget, configured centrally (in the same tier config, keyed by endpoint pattern) rather than hardcoded per-endpoint in gateway code.
Failure Handling
| Failure | Detection | Response |
|---|---|---|
| Redis (shared counter store) unavailable | Sync worker connection errors | Gateway instances continue enforcing against their local counters only, resetting to a conservative per-instance cap; degraded accuracy, not degraded availability |
| Config store unavailable | Local config cache miss on a new client | Fall back to a global default tier (most restrictive reasonable limit) rather than failing the request or allowing unlimited traffic |
| Local counter cache grows unbounded (many distinct clients) | Memory alert on gateway instance | LRU eviction on the local cache — an evicted client’s next request just re-initializes from Redis, costing one network hop instead of zero |
| Clock skew between gateway instances | Sync reconciliation shows implausible deltas | Window boundaries are computed from a shared, coarse-grained time source (synced via NTP, checked periodically) rather than trusting each instance’s raw clock for anything more precise than second-level buckets |
| Sync worker falls behind (backlog) | Sync queue depth alert | Local-only enforcement continues (correct-ish, per above); if backlog persists past a threshold, alert rather than silently drifting indefinitely |
Scaling
Horizontal scaling of the gateway fleet
Because enforcement is local-first with async reconciliation, adding gateway instances doesn’t add proportional load to the shared counter store — each new instance still syncs on the same fixed interval, not per-request, so the shared store’s load scales with the number of distinct clients and the sync frequency, not with total request volume. This is what keeps the design viable well past 500,000 requests/second.
Shared counter store scaling
Redis is sharded by clientId (consistent hashing), so a single hot client’s counter updates land on one shard while the other 99,999 clients’ traffic spreads across the rest — this avoids the scenario where one very active client’s sync traffic becomes a shard hotspot for everyone sharing that shard.
Scaling to more clients
At 100,000 clients the working set is trivially small (a few megabytes, per Capacity Estimation); this design scales to millions of tracked clients without a fundamental rethink, since both local cache size and shared counter storage grow linearly and cheaply per client, and eviction on the local side bounds per-instance memory regardless of total client count.
Geo-distributed gateways
For a gateway deployed across multiple regions, each region’s local counters sync to a regional Redis cluster, and regional clusters cross-replicate counter deltas (not full state) on a longer interval — a client hitting both a US and EU region gets limited against a slightly-more-stale global view, an acceptable widening of the same soft-overshoot trade-off already accepted within a single region.
Centralized counting (every check hits shared Redis)
- Exact, immediately consistent limits
- Hard latency floor of a network round trip on every request
- Redis becomes a scaling ceiling and a single dependency on the hot path
Local-first counting with async reconciliation
- Sub-microsecond hot-path check, scales with gateway fleet size
- Soft overshoot possible — up to roughly instanceCount x limit in a worst case
- More moving parts: sync workers, eviction, staleness to reason about
Recommendation — Local-first for general API gateway limiting, where the goal is fairness and backend protection, not exact billing-grade accounting. Reserve centralized counting for the specific limits tied directly to a hard quota or billing guarantee.
Observability
- Reject rate per client and per endpoint — the primary product-facing signal, both for spotting abusive clients and for spotting a limit that’s misconfigured too aggressively for legitimate usage patterns.
- Local-versus-shared counter drift — sampled comparison between a gateway instance’s local count and the authoritative Redis count for the same client, which is the direct measurement of how much soft overshoot the system is actually producing in practice, not just in the worst-case theoretical bound.
- Sync worker lag and Redis operation latency — the leading indicator that the fallback-to-local-only failure mode is about to engage, ideally caught and alerted on before it does.
- Check latency at the gateway (p50/p99), specifically isolating the rate-limiter’s own contribution to total request latency, since this is a component explicitly budgeted at under 5ms and any regression here is felt on every single request the gateway serves.
Rejected requests get a lightweight structured log line (clientId, endpoint, tier, limitValue, currentEstimate) sampled at a rate that keeps volume manageable even during a sustained abuse event, since a client hammering a limit can otherwise generate log volume proportional to the very traffic the limiter is meant to control.
Trade-offs
The defining trade-off in this design is local-first approximate counting versus centralized exact counting, covered in depth in the Deep Dive and Tradeoff callout above. The whole architecture is built around accepting bounded, small inaccuracy in exchange for keeping the rate limiter’s own overhead negligible relative to the traffic it protects — a rate limiter that adds more latency than it saves in backend protection has failed at its actual job, even if its counting is perfectly precise.
The second trade-off is running two different algorithms (token bucket for burst, sliding window counter for sustained) rather than one uniform approach everywhere. This costs implementation and operational complexity — two code paths to test, reason about, and debug — but a single algorithm forced to serve both burst-tolerance and sustained-rate-smoothing ends up doing both jobs worse than either specialized approach does its one job, which is a bad trade for a component that’s supposed to be simple to reason about under incident pressure.
Final Architecture
The finished system checks two independently-tuned algorithms on every request entirely in local memory, keeps the shared source of truth updated asynchronously rather than synchronously, and degrades to local-only enforcement rather than failing the gateway open or closed when that shared store has a problem. The result is a rate limiter whose own cost is close to invisible against the traffic it’s protecting, at the price of a small, bounded, and deliberately accepted amount of overshoot.