Overview
Curated: · Written: · Reviewed:
Key takeaways
- Scale the measured bottleneck, not the component that is easiest to duplicate.
- Vertical and horizontal scaling are complementary. Both have ceilings, costs, and failure modes.
- Stateless application instances improve interchangeability, but state still exists and its owner must scale.
- Load balancing distributes work; it does not create backend capacity or guarantee even utilization.
- Autoscaling is delayed feedback. Preserve headroom and protect the warm-up gap with admission control.
- Design overload behavior before launch: bounded queues, deadlines, retry budgets, load shedding, and graceful degradation.
- Validate capacity and failover with production-shaped load, percentiles, saturation, and dependency telemetry.
1. Define the objective
Scalability is the ability to meet a changing workload while preserving an explicit service objective at an acceptable cost. State the workload in useful units—requests per second, concurrent sessions, bytes per second, queued jobs, or writes per partition—and define latency percentiles, error rate, freshness, and availability targets. “Handle millions of users” is not a capacity model.
Latency is time per operation; throughput is completed work per unit time. They can trade off: batching improves throughput by amortizing overhead while adding queueing delay. Measure the distribution, not only the mean. A healthy average can hide an unacceptable p99.
2. Find the bottleneck
A bottleneck is workload-dependent and can move. Increase representative load and correlate traffic with latency, errors, and saturation at every mandatory dependency. CPU, memory, connection pools, locks, queue depth, storage IOPS, network bandwidth, downstream quotas, and database contention can each become the active constraint. Optimize or scale the first limiting resource, then repeat the test because the next constraint may appear.
3. Vertical and horizontal scaling
Vertical scaling gives one node more CPU, memory, storage, or bandwidth. It is often the simplest and fastest response, but instance sizes, single-node throughput, maintenance, and failure-domain requirements impose ceilings.
Horizontal scaling adds independently schedulable workers or replicas. It can increase aggregate capacity and support redundancy only when work can be partitioned, shared dependencies have headroom, state and coordination are designed, and replicas span the intended failure domains. It is not unbounded: coordination, skew, hot partitions, fan-out, network limits, and serial work eventually dominate.
Use both pragmatically. Choose an efficient node size, add nodes as demand grows, and retain enough spare capacity to survive the failures and deployments named in the reliability target.
4. Statelessness and state ownership
A stateless request handler does not rely on process-local session state surviving between requests. Any eligible instance can handle the next request. This simplifies replacement and balancing, but it relocates state to cookies, databases, caches, or durable services. Those systems become shared dependencies with their own latency, capacity, consistency, and availability contracts.
Sticky routing can preserve local state or locality, but it creates skew and makes failover harder; it is a deliberate tradeoff, not proof that horizontal scaling is impossible.
5. Load balancing and health
Round robin works for roughly equal work. Least-connections is only a proxy for load: one connection may carry very different work, especially with multiplexing. Weighted policies encode known capacity differences; hashing improves affinity but can create hot keys.
Separate process liveness from traffic readiness. A readiness check should prove the instance can accept its intended traffic without turning every shared-dependency incident into removal of the entire fleet. Use thresholds and hysteresis, observe probe traffic separately, and test startup, overload, dependency failure, and recovery.
6. Autoscaling and capacity
Reactive autoscaling observes a signal, decides, provisions, starts, and warms capacity. It therefore lags sudden demand. CPU can be a poor signal for I/O-bound work; concurrency, queue age, request rate, or a workload-specific saturation signal may lead better. Scheduled scaling helps predictable peaks. Mature designs combine both with minimum capacity, bounded fallback, and scale-in protection.
Plan for peak workload plus the capacity lost in the chosen failure scenario—not a universal percentage. Validate the plan with production-shaped data and user journeys, step and spike tests, soak tests, dependency limits, and failure injection.
7. Overload control
An unbounded queue converts overload into memory growth and stale work. Use bounded queues, admission control, end-to-end deadlines, cancellation, and backpressure. When offered load exceeds useful capacity, reject cheaply before the service enters a latency and retry spiral.
Rate limiting enforces a policy per identity or resource; load shedding protects current service health. HTTP 429 can express a client rate limit, while 503 is commonly used for temporary service overload. Clients should retry only safe transient failures, with a limit, exponential backoff, random jitter, and respect for Retry-After where supplied. Retrying at several layers multiplies attempts.
Graceful degradation serves a cheaper result or disables noncritical work. Exercise this mode regularly; an emergency path that is never tested is not a reliability mechanism.
8. Availability and parallelism
Multiplying component availability is valid only when every component is mandatory, events are independent, and the measurement windows match. Shared failure modes, fallbacks, redundancy, partial operation, and correlated maintenance change the model.
Amdahl's Law provides a bound for a fixed task: if fraction s is serial, ideal speedup with N workers is 1 / (s + (1-s)/N). Real services add coordination and queueing overhead, so measure rather than treating the formula as a capacity forecast.
Interview framework
State the workload and SLO, draw the mandatory dependency path, calculate normal and failure-mode load, identify the current bottleneck, choose a scaling mechanism, protect overload and retries, and name the load test and telemetry that would falsify the design.
Worked example: finding the bottleneck before scaling anything
A service is offered 3,000 requests/second. It completes 2,667 of them, sheds the rest at the application edge, and the ones that complete have a p99 of 1,850 ms against a 500 ms objective. Keeping offered load and goodput separate matters here — this section defines throughput as completed work, and 11% of what arrives is never served. The instinct is to add application instances. Measuring each tier says why that would spend money and change nothing:
| Tier | Utilization | Time spent in this tier | Note |
|---|---|---|---|
| Load balancer | 12% | 3 ms | |
| Application (24 instances) | 34% | 61 ms | |
| Database connection pool | saturated | ~1,500 ms waiting for a connection | 40 connections, wait queue bounded at 4,000 |
| Database query execution | 41% CPU | 15 ms mean per query | |
| Cache | 22% | 8 ms |
Percentiles do not add, so these columns are not meant to reconstruct the 1,850 ms headline. The point is that one tier dominates every other: the pool wait is about 25 times the application tier and about 100 times query execution. It is also not the tier anybody was about to scale.
The pool number is a capacity statement, not a queueing approximation:
pool capacity = connections / mean query time = 40 / 0.015 s = 2,667 req/s (goodput)
offered load = 3,000 req/s
deficit = 333 req/s = 11% of offered, shed at the edge once the queue is full
wait, admitted = 4,000 waiting / 2,667 completed per second = 1.50 s
Two things follow, and the second is the one that gets missed. The wait is set by how deep the queue is allowed to get, not by how much the pool missed by: at 4,000 waiters it is 1.5 seconds, at 400 it would be 150 ms. And the queue only stays at a finite 4,000 because the 333/s that cannot be admitted are shed — a queue that admits everyone has no steady state at all, and the 1.5 s figure would be meaningless.
So the p99 is not the worst number on this page. The service is failing 11% of what it is asked to do, and no percentile computed over the requests that completed will ever show that. Goodput and latency have to be read together.
Adding 24 more application instances leaves that tier 17% utilized and changes neither number, because none of those instances can get a connection either. Raising the pool from 40 to 60 gives 60 / 0.015 = 4,000 req/s of capacity against 3,000 offered, so the queue drains, shedding stops, goodput reaches the full 3,000/s, and the pool settles at 75% utilization.
The database does not become the next bottleneck, which is the part worth getting right. Steady-state concurrency is arrival rate times service time, not pool size: at 3,000 req/s and 15 ms per query, 45 queries are in flight whether the pool holds 60 connections or 600. Database CPU therefore rises from 41% to about 46% — idle connection slots are not extra work. CPU becomes the constraint only if offered load keeps growing, and the same measurement has to be repeated then rather than predicted now.
