Scalability
SystemsHigh priority~20 minMeeting more load by adding capacity rather than rewriting — what the claim actually asserts, the two dimensions capacity comes in, and why the tenth machine returns less than the second.
Definition
Scalability is the property that more load can be met by adding capacity rather than by redesigning, at a cost that grows no faster than the load it absorbs. The payoff is optionality: a traffic forecast becomes a purchasing decision instead of a rewrite.
You pay for it in coordination. Once capacity is more than one machine, a request may need state that lives somewhere else, and that hop is latency, a new failure surface, and a consistency question you previously did not have to answer.
A claim of scalability is a claim about a curve, so it means nothing until both axes are named: the load quantity that grows, and the resource cost that answers it.
Doubling read traffic costs one more replica and no code change; doubling writes costs a resharding names both, and someone can check it — including where the curve ends, because every design scales from here to there and then stops.
When to use
The cue is a prompt that names a growth number instead of a feature — for 10 million users, during Black Friday, it works today but falls over at peak.
The first move is not a diagram. Ask which number grows: read QPS, write QPS, stored bytes, or fan-out per write (The 4-step interview framework owns getting those numbers out of a vague prompt).
When load is flat and the complaint is one slow request, this is the wrong page — that is latency, not throughput.
Techniques
Capacity comes in two dimensions, plus a third move that avoids needing it.
Vertical scaling (scale up) buys a bigger machine. Nothing about the software changes and no request crosses a network it did not cross before, which makes it the cheapest option in engineering time. The ceiling is hard, and far higher than most candidates assume. The counterweight is price and downtime: the largest instances cost disproportionately more per unit of capacity, and the upgrade is usually a restart.
Horizontal scaling (scale out) puts more machines behind something that spreads work across them. There is no ceiling in principle, and AWS's Well-Architected reliability principles recommend it for a second reason: replacing one large resource with several small ones reduces the impact of a single failure. The precondition for the service tier is statelessness — an app instance that remembers anything between requests cannot be one of N. The data tier scales out the other way: it keeps the state and partitions it instead (sharding).
Scaling by doing less work is the one candidates forget: cache, denormalise, precompute, move work off the request path. It buys headroom in staleness rather than in machines.
Related concepts
This page owns the three capacity strategies and the shape of the curve they buy. Everything they are made of is elsewhere: which dimension, when is vertical vs horizontal scaling; the precondition is stateless services behind load balancers.
The per-tier mechanics fill the rest of this track — caching for repeated reads, read replicas and sharding for the database, SQL vs NoSQL for engines built to partition from the start, and message queues to push work off the request path.
Worked examples
Stack Overflow, February 2016 — nine web servers
Nick Craver's write-up of Stack Overflow's architecture publishes a full day of production counters for 9 February 2016: 209,420,973 HTTP requests, of which 66,294,789 were page loads, served by a web tier of nine primary servers plus two for staging. Every figure here is that day's measurement.
The arithmetic below is the point, and Craver states its conclusion more bluntly than the division does: they were down to needing only one web server, and had unintentionally tested that, successfully, a few times.
| Quantity | Value | Derivation |
|---|---|---|
| Largest single EC2 instance AWS documents | 1,920 vCPUs · 32 TiB RAM · 200 Gbps | AWS EC2 U7i product page, size u7inh-32tb.480xlarge — vendor-documented, not measured by us. |
| Stack Overflow HTTP requests, 9 Feb 2016 | 209,420,973 / day | Measured counter published in the 2016 architecture post. |
| → site-wide request rate | ≈ 2,424 req/s | 209,420,973 ÷ 86,400 seconds in a day. |
| → per primary web server | ≈ 269 req/s | 2,424 ÷ 9 primary web servers (numbered 01–09 in the post). |
| Page loads, same day | 66,294,789 / day ≈ 767/s | 66,294,789 ÷ 86,400. Under a third of requests render a page; the rest are lighter. |
| Average question page render | 22.71 ms | Measured across 49,180,275 renders that day; the home page averaged 11.80 ms. |
What makes those numbers possible is the shape of the workload, not a trick. Reads dominate, the hot set is cacheable — their two Redis nodes ran below 2% CPU serving it — and rendering is cheap and measured. Behind that sits vertical scaling doing the heavy lifting: two SQL Server clusters at 384 GB and 768 GB of RAM per box, on PCIe SSDs.
The nine servers are deliberate over-provisioning, and Craver names the reasons: rolling builds, headroom, redundancy. Nine boxes for a one-box load is a different purchase from capacity — it buys the ability to lose several and to deploy without a maintenance window (Availability).
Where this shape stops is where the rest of the track begins: writes that fan out to many readers, a second region, and a working set that outgrows RAM.
Tradeoffs
| What it costs | When it bites | |
|---|---|---|
| The second machine | State an instance kept in memory moves behind a network hop, or is given up | The first request routed to a server other than the one holding its session, upload or cached entry |
| Coordination between nodes | Throughput stops being linear in N and eventually reverses | Write-heavy workloads where nodes must agree — past the peak, each node added returns less than nothing |
| Scaling out early | Deploys, monitoring, partial failure and a distributed bug class, bought for capacity one box already had | A sharded cluster proposed for a few hundred requests a second per server |
| Doing less work instead | Staleness, and an invalidation problem that outlives whoever added the cache | The first read that has to reflect a write immediately — the boundary CAP theorem and consistency models draw |
The second row has a model behind it. Neil Gunther's Universal Scalability Law writes throughput at load N as X(N) = γN / (1 + α(N−1) + βN(N−1)). Its three coefficients are γ for ideal linear speedup, α for contention (queueing on a shared resource), and β for coherency (the delay while distributed copies of data agree). With β = 0 and γ = 1 it reduces to Amdahl's law, whose ceiling is 1/α: added nodes get you asymptotically nowhere, but never backwards.
β > 0 is worse than a ceiling. The coherency term grows quadratically, so throughput peaks at Nmax = √((1−α)/β) and decreases beyond it. Take α = 0.05 and β = 0.005 — coefficients assumed here to make the shape visible, not measured from any system. Nmax = √(0.95 ÷ 0.005) ≈ 14 nodes, worth 5.5× one node's throughput. Thirty-two nodes return 4.3×: more hardware, less work done.
Which is why scale out is an answer only once you can say what the added nodes have to agree about.
Things to look out for
- Scaling the tier that wasn't the bottleneck. Adding app servers to a saturated database makes the database worse — more connections, more lock contention, the same disk. Name the saturated resource before adding anything.
- The shared thing behind the tier you scaled. A connection pool, a licence server, one writable primary, a single cache node: N instances in front of one of something is still one of that something (single points of failure).
- Sizing on the average when the prompt named a peak. A daily mean says nothing about the Monday-morning spike, and hardware is bought for the spike. Quote both, and quote response time at a percentile (latency vs throughput).
- Cost climbing faster than load. A design that absorbs 10× the traffic for 30× the bill has not scaled, it has deferred. Carry a per-request cost beside the capacity plan.
In the interview
- You've added a second app server — what breaks? The depth answer is a list of state: sessions in memory, in-process caches, uploads on local disk, scheduled jobs that now fire twice, anything that assumed there was one of it. Then name the load balancer you just made critical (single points of failure).
- How far does one machine get you? Refuse the abstract version. Quote the documented ceiling — one EC2 instance is 1,920 vCPUs and 32 TiB — then the independent point that Stack Overflow's own engineers said nine web servers was eight more than they needed, on 2016 hardware. Arguing for the simpler architecture with evidence lands harder than proposing a cluster.
- Why doesn't doubling the nodes double the throughput? Name both costs — contention for shared resources, coherency between copies — then the retrograde region past Nmax, where added nodes subtract throughput.
- Capacity is always added to a tier, so the follow-up is which one. Say which saturates first and why you believe it — a number, not a hunch — before naming any mechanism.