2.4 Scalability
For a new product with few users, the overriding engineering goal is keeping the system simple and flexible so you can adapt as you learn what customers need.
Scalability = a system's ability to cope with increased load.
4.1 The "you're not Google" caveat, handled properly
For a new product with few users, the overriding engineering goal is keeping the system simple and flexible so you can adapt as you learn what customers need. Worrying about hypothetical future scale is counterproductive:
- Best case: wasted effort and premature optimization
- Worst case: it locks you into an inflexible design and makes the application harder to evolve
4.2 Scalability is not a label
It is meaningless to say "X is scalable" or "Y doesn't scale." The real questions:
- If the system grows in a particular way, what are our options for coping?
- How can we add computing resources to handle the additional load?
- Based on current growth projections, when will we hit the limits of our current architecture?
4.3 Understanding load
You must first describe current load before you can ask "what if it doubles?"
Usually a throughput measure: requests/second, GB of new data per day, checkouts per hour. Sometimes the peak of a variable quantity (simultaneously online users).
Other statistical characteristics that change the access pattern and hence the scalability requirement:
- read:write ratio
- cache hit rate
- number of data items per user (followers, in the case study)
- whether the average case matters, or whether the bottleneck is dominated by a small number of extreme cases
Then ask the growth question in two symmetrical forms:
- Fix resources, increase load → how is performance affected?
- Fix performance, increase load → how much must resources increase?
Goal: stay within SLA while minimizing cost.
Linear scalability: doubling resources handles twice the load at the same performance — considered good. Occasionally you do better than linear (economies of scale, better distribution of peak load). Much more likely, cost grows faster than linearly — e.g., with a lot of data, processing a single write may involve more work than with a small amount of data, even for the same request size.
4.4 Shared-memory vs shared-disk vs shared-nothing
| Shared-memory — scale up | Shared-disk | Shared-nothing — scale out |
|---|---|---|
| Many CPUs share one pool of RAM and disks; threads and processes share memory directly. | Independent machines with their own CPU and RAM, sharing disks over a fast network (NAS / SAN). | Each node has its own CPU, RAM and disk. Coordination happens in software only, over the network. |
| Cost grows faster than linearly, and internal bottlenecks mean 2× the hardware does not carry 2× the load. | Contention and locking overhead limit how far it scales. | Can scale linearly on commodity hardware, elastic in both directions, and tolerates the loss of a datacenter. |
| — | — | Requires explicit sharding, and brings all the complexity of distributed systems with it. |
| Architecture | Also called | Reality |
|---|---|---|
| Shared-memory | vertical scaling, scaling up | Individual CPU cores aren't getting much faster, but you can buy more cores/RAM/disk. Parallelism via threads sharing RAM. Cost grows faster than linearly — a machine with 2× resources costs significantly more than 2×, and because of bottlenecks it likely can't actually handle 2× the load |
| Shared-disk | — | Multiple machines with independent CPU/RAM, data on a shared disk array over a fast network (NAS / SAN). Traditionally used for on-premises data warehousing. Contention and locking overhead limit scalability |
| Shared-nothing | horizontal scaling, scaling out | Each node has own CPU, RAM, disk; all coordination in software over a conventional network. Potential for linear scaling, best price/performance hardware (especially in cloud), easy resource adjustment, fault tolerance across datacenters and regions. Costs: explicit sharding (Ch 7) and all the complexity of distributed systems (Ch 9) |
The cloud-native variant. Some cloud-native databases use separate services for storage and transaction execution, with multiple compute nodes sharing one storage service. Superficially like shared-disk — but it avoids the old scalability problems because the storage service exposes a specialized API designed for the database's specific needs, rather than a filesystem (NAS) or block device (SAN) abstraction. That API-level specialization is the whole trick.
4.5 Principles for scalability
There is no generic, one-size-fits-all scalable architecture ("magic scaling sauce").
Illustration: a system for 100,000 requests/s × 1 kB looks very different from one for 3 requests/minute × 2 GB — even though both move ~100 MB/second.
An architecture appropriate for one level of load is unlikely to cope with 10× that load. On a fast-growing service you will probably rethink your architecture on every order-of-magnitude load increase. It is usually not worth planning more than one order of magnitude ahead.
Principle 1 — decompose. Break the system into smaller components that operate largely independently. This underlies microservices, sharding (Ch 7), stream processing (Ch 12), and shared-nothing. The challenge is knowing where to draw the line between what belongs together and what belongs apart.
Principle 2 — don't over-complicate.
- If a single-machine database will do the job, prefer it to a complicated distributed setup.
- Autoscaling is cool, but if your load is fairly predictable, a manually scaled system may have fewer operational surprises.
- A system with 5 services is simpler than one with 50.
- Good architectures usually involve a pragmatic mixture of approaches.