Most scaling advice is a list of techniques: add a cache, add replicas, shard the database, put a queue in front of it. Each is sound. The trouble is order and fit. A cache does nothing for a write-bound system. Sharding a database whose real problem is one hot row just gives you a distributed hot row. Ten more web servers in front of a connection-limited database make things worse, because each one opens its own pool.

This article is about the decision method that sits above the techniques. First, identify which resource is saturated. Then work out how much of it you need using arithmetic you can do on a whiteboard. Finally, pick the cheapest move that relieves that resource and measure again. The individual techniques have their own pages on this site; here they appear as rungs on a ladder, each with the condition that justifies climbing to it.

Advertisement

Scaling starts with naming the saturated resource

A system stops scaling when one resource reaches its limit and requests start queueing for it. That resource is almost always one of a short list: CPU on the application tier, memory (including garbage-collection pressure), database connections, database CPU or lock contention, disk IOPS, network bandwidth, or a downstream dependency with its own quota. Symptoms mislead. High latency on the API may come from the database; high database CPU may come from a missing index rather than from load.

The discipline is the USE method applied per resource: utilisation, saturation (work waiting), and errors. A resource at 95% utilisation with a growing queue is the bottleneck. A resource at 95% with no queue is merely busy. Collect those three numbers for every tier before choosing a pattern, and write down which one crosses its limit first as load rises. That resource, and only that one, is what the next move should relieve.

It also helps to separate throughput problems from latency problems. If requests per second plateau while the servers are busy, capacity is the issue. If throughput is fine but the 99th percentile is bad, the issue is usually queueing near saturation, tail-heavy dependencies, or garbage-collection pauses. The second kind is often fixed by running at lower utilisation.

The scaling ladder: name the saturated resource, then take the cheapest move that relieves it1. Measurewhich resource is at its limit?2. Verticalbigger box, zero code change3. Stateless tierstate out, load balancer inREAD PATHWRITE PATH4a. Cachehit ratio removes reads from the database5a. Read replicaslag-tolerant reads fan out4b. Async and batchingqueue absorbs bursts, batches cut round trips5b. Shardsplit write ownership by key6. Protect at the limitbackpressure, load shedding and autoscaling keep overload from becoming collapseEvery rung is re-measuredLittle's law gives the concurrency you must hold; the USL curve tells you when adding nodes stops paying
The ladder. Measurement comes first and returns after every rung. Vertical scale and a stateless tier are nearly always the first moves; after that the read path and write path diverge, and protection at the limit applies to both.

The capacity arithmetic: Little's law and queueing

Little's law says the average number of requests in a system equals arrival rate times time in system: L = λ × W. At 2,000 requests per second with a 50 ms average service time, 100 requests are in flight on average. That number sizes everything that holds a request: worker threads, event-loop slots, database connections, and the concurrency setting on a serverless platform. If each application instance can hold 20 concurrent requests, you need at least 5 instances before any headroom.

Queueing theory explains why you cannot run at 100%. For a single server with random arrivals, the time in system is roughly S / (1 − ρ), where S is service time and ρ is utilisation. At 50% utilisation a 10 ms request takes about 20 ms. At 80% it takes 50 ms. At 95% it takes 200 ms. Latency grows without bound as utilisation approaches one, and the tail gets there first. That is why capacity targets of 60 to 70% are normal for latency-sensitive tiers, and why autoscalers target a utilisation below the knee rather than at the limit.

Adding nodes does not raise throughput linearly either. Neil Gunther's Universal Scalability Law models throughput at N nodes as X(N) = λN / (1 + α(N − 1) + βN(N − 1)). The α term is contention: the fraction of work that is serialised, such as a shared lock or a single leader. The β term is coherency: the cost of nodes keeping each other up to date, such as cache invalidation or cross-shard coordination. With β above zero, throughput peaks and then falls. Fitting the curve to a load test tells you where horizontal scale stops paying:

# Fit the Universal Scalability Law to a load test and find the useful node count.
# X(N) = lam * N / (1 + a*(N-1) + b*N*(N-1))
#   a = contention (serialised fraction), b = coherency (cross-node chatter)
import numpy as np
from scipy.optimize import curve_fit

nodes = np.array([1, 2, 4, 8, 16, 24])
rps   = np.array([950, 1850, 3500, 6200, 9800, 11100])   # measured throughput

def usl(n, lam, a, b):
    return lam * n / (1 + a * (n - 1) + b * n * (n - 1))

(lam, a, b), _ = curve_fit(usl, nodes, rps, p0=[1000, 0.02, 0.0001], bounds=(0, np.inf))
n_peak = np.sqrt((1 - a) / b) if b > 0 else float("inf")
print(f"lambda={lam:.0f} rps/node  contention={a:.3f}  coherency={b:.5f}")
print(f"throughput peaks near N={n_peak:.0f}; beyond that, adding nodes lowers throughput")

If the fitted α is large, look for the serialised step, since no amount of hardware fixes it. If β dominates, reduce cross-node chatter: partition so that nodes stop sharing state.

Advertisement

Rung one: scale up before you scale out

Vertical scaling means a larger machine: more cores, more memory, faster disks. It needs no code change, keeps a single copy of the data, and avoids every distributed-systems problem. For a database in particular, a bigger primary is almost always cheaper in engineering time than sharding.

Its limits are real, though. Price rises faster than capacity at the top of the range. One machine is one failure domain, so vertical scale must be paired with a standby. Resizing usually means a restart or failover. And software can stop using the extra hardware: a single-threaded process, a global lock or a connection limit caps the gain regardless of core count. Scale up while the price and risk are acceptable, and use the time it buys to prepare the next rung.

Rung two: a stateless application tier behind a load balancer

Horizontal scaling of the application tier depends on one property: any instance can serve any request. That means no user state in process memory, no local files that other instances need, and no scheduled jobs that assume one copy. Session data moves to a shared store or into a signed token; uploads go to object storage; periodic jobs take a lock or move to a scheduler.

# Before: session in process memory -> node affinity, lost carts on every deploy
@app.post("/cart/add")
def add(item_id):
    SESSIONS[request.cookies["sid"]]["cart"].append(item_id)

# After: session in a shared store with a TTL -> any node can serve any request
@app.post("/cart/add")
def add(item_id):
    sid = request.cookies["sid"]
    key = f"cart:{sid}"
    pipe = redis.pipeline()
    pipe.rpush(key, item_id)
    pipe.expire(key, 60 * 60 * 24)   # abandoned carts expire instead of leaking
    pipe.execute()

Once instances are interchangeable, a load balancer spreads requests and health checks remove failed instances. Prefer least-outstanding-requests over round robin when request cost varies widely. Avoid sticky sessions as a crutch: they reintroduce affinity, unbalance load after scale-out, and turn every deploy into lost state. The main trap at this rung is downstream: every instance opens its own database pool, so 40 instances with 20 connections each ask the database for 800 connections. Size pools from Little's law, not from defaults, and put a connection pooler in front of the database when instance counts grow.

The read path: cache first, then replicas

Most systems read far more than they write, so the read path usually saturates the database first. The cheapest relief is a cache. A 90% hit ratio removes nine of ten reads, so the database sees a tenth of the load. The costs are staleness, invalidation logic, and a new failure mode when the cache is cold or down. The patterns (cache-aside, write-through, TTLs, stampede protection) are covered in caching in system design.

Read replicas come next for reads that miss the cache or cannot be cached, such as search and reporting queries. Replicas copy the primary asynchronously, so a read immediately after a write may not see it. Routing must therefore distinguish lag-tolerant reads from read-your-writes reads, which is the subject of read replica routing. Replicas do nothing for writes. Every replica applies every write, so write capacity is still that of one primary.

The write path: absorb, batch, then shard

When writes are the bottleneck, first ask whether each write must complete inside the request. Many writes do not: sending an email, updating a search index, recording an analytics event. Moving them to a queue with workers flattens bursts, since the queue absorbs peaks and workers drain at a steady rate. It also lets you batch: one insert of 500 rows is far cheaper than 500 single-row inserts, because per-statement overhead and commit cost dominate small writes.

When the synchronous write rate itself exceeds one primary, the remaining move is to split ownership of the data by key so that each shard takes a fraction of the writes. Sharding brings shard-key choice, cross-shard queries, rebalancing and hot keys, all covered in database sharding architecture. It is the most expensive rung, and it is mostly irreversible, because application code comes to depend on the key. Climb to it only when the arithmetic shows that one primary, even the largest available, cannot hold the write rate at target utilisation.

Protecting the system at its limit

Every system has a limit, and demand will sometimes exceed it. Without protection, overload turns into collapse: queues grow, latency exceeds client timeouts, clients retry, and the retries add load that no one is waiting for. Three mechanisms keep overload bounded. Backpressure slows producers when consumers fall behind, as described in backpressure. Load shedding rejects excess work early and cheaply, preferring low-priority traffic, as covered in load shedding. Autoscaling adds capacity, but only on a timescale of tens of seconds to minutes, so it cannot absorb a spike on its own. Pair it with shedding and with retry budgets on clients.

Worked example: a catalogue API from 1,500 to 12,000 requests per second

Start with a catalogue API serving 1,500 requests per second at a 250 ms p99 target. It runs on two 16-vCPU instances against one database primary. Profiling shows 6 ms of application CPU per request and three database queries per request, 95% of them reads.

Application tier: 1,500 × 6 ms = 9 CPU-seconds per second, so about 9 cores busy out of 32, or 28%. Database: 4,500 queries per second. Its CPU sits at 70% and p99 query latency is climbing, so the database is the saturated resource. More application instances would not help.

Step 1 is a cache for product detail reads. With an 85% hit ratio on the 4,275 read queries per second, about 640 reads per second still reach the database, plus 225 writes. Database CPU drops to about 15%. Step 2 comes at 6,000 requests per second, when the application tier reaches 36 busy cores. It moves to eight stateless instances. Little's law at 6,000 requests per second and 40 ms mean latency gives 240 in flight, so each instance needs room for about 30 concurrent requests plus headroom. Connection pools are sized at 10 per instance behind a pooler, not at a default of 50.

Step 3 arrives at 12,000 requests per second. Cache misses and search queries now load the primary again, so two replicas take search and listing reads. Step 4 is not needed: 1,800 writes per second, with batching, is still within one large primary. This is the common outcome. Most systems never need to shard, because caching and replicas relieve the read path and batching relieves the write path.

Failure modes

  • Scaling the wrong tier: adding application instances when the database is saturated raises connection count and makes latency worse.
  • Connection explosion: pool size multiplied by instance count exceeds the database limit during scale-out or a deploy surge.
  • Cold cache after a deploy or failover: hit ratio drops to zero and the full read load lands on the database at once; warm caches or ramp traffic gradually.
  • Retry storms: timeouts trigger retries that multiply load exactly when capacity is shortest; cap retries with budgets and jittered backoff.
  • Hot keys: one celebrity item or tenant concentrates load on one cache node or shard regardless of how many exist.
  • Autoscaler lag: scale-out takes minutes while the spike takes seconds; without shedding, the spike becomes an outage before new capacity is ready.

Trade-offs at a glance

MoveRelievesCostsReversible?
Vertical scaleAny single-node limitPrice at the top end, one failure domainYes
Stateless tier and load balancerApplication CPU and memoryState externalisation, more DB connectionsYes
CacheRead load on the databaseStaleness, invalidation, cold-start riskYes
Read replicasUncacheable readsReplication lag, routing logicYes
Async and batchingWrite bursts, per-write overheadEventual completion, queue operationsMostly
ShardingWrite rate and data size beyond one primaryCross-shard queries, rebalancing, hot keysRarely

What to do next

  1. For each tier, record utilisation, saturation and errors at current peak, and write down which resource will hit its limit first.
  2. Compute in-flight requests with Little's law and check that threads, pools and platform concurrency settings match it with headroom.
  3. Run a stepped load test at 1, 2, 4 and 8 instances and fit the USL; find the node count where throughput stops rising.
  4. Remove in-process state from the application tier and size database pools from instance count times pool size.
  5. Measure the read/write split and cache hit ratio before considering replicas or shards.
  6. Add load shedding and client retry budgets before relying on autoscaling; see cloud autoscaling for its reaction times.
Key takeaway: Scaling is a loop, not a menu. Name the saturated resource, size it with Little's law and a utilisation target below the queueing knee, and take the cheapest move that relieves that resource: scale up, make the tier stateless, cache, replicate reads, move writes off the request path, and shard only when the write arithmetic demands it. Measure again after every move, and protect the limit with shedding and backpressure, because demand will eventually exceed whatever you build.