Design an API that enforces a global per-client request limit correctly even after the serving tier is scaled to many instances, where naive per-instance counters silently fail.
You are running an API that enforces a rate limit per client, for example N requests per minute. It works perfectly on a single server. Then traffic grows and you scale the serving tier horizontally to many application servers behind a load balancer, and the rate limit quietly stops working : clients routinely exceed N, and no single server ever sees them do it.
The reason is the classic distributed-systems trap. When each application server keeps its own in-memory counter, every instance only counts the fraction of a client's traffic that the load balancer happened to route to it. With the traffic spread across many servers, each instance independently believes the client is under the limit, while the true aggregate sails well past N. The limit is only correct if every server consults a single shared view of a client's usage, rather than its own local tally.
Design the architecture so the rate limit is enforced globally across the whole serving tier. The limit state must live in one shared place that every application server consults, not be replicated per instance. Lay out the topology (client, load balancer, a horizontally scaled serving tier, the shared limiting mechanism, and the datastore), then document the API and the trade-offs of centralizing limit state, such as the added dependency on the shared limiter and the latency of consulting it on every request.