System evolution

Step a product from a single VM to a sharded, replicated architecture - splitting the DB, measuring before scaling, load-balancing, caching, moving slow work to queues, then replicas and shards. Each stage shows the bottleneck it fixes and the trade-off it brings.

9 min read

No one designs for millions of users on day one, and nobody should. Real architectures evolve one bottleneck at a time. The skill a design conversation is testing is almost never “can you draw the big diagram”. It is knowing which piece the system needs next, and why.

The final diagram is the wrong answer to almost every question

Asked to design a system, most people reach straight for the last stage: shards, replicas, queues, several regions. It is the most expensive possible way to serve a few thousand users. Every component buys capacity and costs something in correctness, latency and things that can break.

1A single VMneeded
2Split the database offneeded
3See what is actually happeningneeded
4More than one app serverState can no longer live in the process; deploys, sticky sessions and config drift get harder.from ~20K
5Cache and CDNCache invalidation is now your problem: stale reads, and a stampede whenever the cache goes cold.from ~100K
6Move slow work off the requestEventual consistency, retries, dead-letter queues, and one more system to monitor.from ~200K
7Go multi-regionData must stay in sync across regions; async replication means writes are not instantly visible everywhere.from ~2M
8Read replicas, then shardsReplication lag makes fresh writes look missing, and sharding ends cross-shard joins and transactions.from ~5M
At 5K daily users you need 3 of 8 stages. Building the rest now means signing up for the 5 amber problems above, for capacity you do not use.
5K
At 5,000 users the first three stages cover it. Everything else is cost without a reason yet.

At 5,000 users none of that extra capacity is needed and all of its cost is real: connection pools, cache invalidation, replication lag, partial failures, and the tracing you need to understand any of it. The right answer is one web server, one database, and monitoring good enough to see the first bottleneck coming.

Walk the stages and watch each piece earn its place

Each stage below exists because the previous one ran out of a specific resource. Step through them and read the three notes each time: what broke, what was added, and what that addition now makes harder.

5Cache and CDN
about hundreds of thousands
metrics, logs and traces over everything below
Clients
Users
Edge
CDNLoad balancer
App
app 1app 2app 3
Cache
Redis
Data
Database
What broke
The same reads ran over and over, and every image and script came from your origin.
What we added
A shared Redis cache in front of the database for hot data, and a CDN at the edge for static assets.
What it costs
Cache invalidation is now your problem: stale reads, and a stampede whenever the cache goes cold.
Each stage exists because the previous one ran out of something specific. Step back to the single VM, or on to shards; new pieces are outlined.

Read only the red notes from start to end and you see the real story: each bottleneck is created by the fix before it. Note also the order of the database fixes. A read-heavy database is the easiest problem on the list, and sharding is the most expensive fix available. Cache first, replicas second, and shard only when writes, which replicas do not help, are what is saturating you. Bigger hardware is a legitimate option too, and it buys real time.

Underneath it all are a handful of principles: keep the web tier stateless, build redundancy at every tier, cache what you can, serve static assets from a CDN, run in more than one data center, shard the data tier when you must, and monitor everything.

The skill is knowing which stage you are at

Every box in that diagram was added to relieve pressure someone observed. So the prerequisite for the whole playbook is being able to see the pressure, and that is the one component nobody draws on the whiteboard.

What are you seeing?
Add a cache, then read replicas. Do not shard yet.That is stage 5, “Cache and CDN”.
Each signal points at one move. Without the monitoring, every one of them is a guess.

Without monitoring, every decision after stage one is a guess. You cannot tell whether the database is read-bound or write-bound, whether p99 is climbing while the mean looks fine, or whether the thing you just added helped. Teams that scale well are not the ones who guessed the right architecture. They are the ones who could see which resource ran out first.

One box and a database, monitored well enough to see the wall coming, is a better answer than a distributed system serving a thousand people.

The short version

  • Every component adds capacity and a permanent cost; add it when the pressure is real.
  • Each stage relieves the bottleneck the previous fix created.
  • For a busy database: cache, then replicas, and shard only for writes.
  • Monitoring comes first, because it tells you which stage you are at.