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.
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.
- 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.
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.
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.