learning

Phase 8: System Design

1370 words7 min read
Phase 8: System Design
Authors

Welcome to System Design, the pinnacle of software engineering.

When you build a side project, your architecture is simple: a frontend, a web server, and a database, all running on a single machine. But what happens when your app gets featured on the news and traffic spikes from 10 users a minute to 100,000 users a second?

A single machine will burst into metaphorical flames. It will run out of RAM, the CPU will hit 100%, and the database will lock up. System Design is the discipline of architecting systems that can handle massive scale without crashing. It is about anticipating bottlenecks, eliminating single points of failure, and understanding the deep tradeoffs in distributed systems.


1. The Scaling Dilemma

When your server hits its limit, you have two choices.

Vertical Scaling (Scale Up)

Vertical scaling means buying a bigger, faster computer. You upgrade your server from 8GB of RAM to 256GB of RAM. You upgrade to a 64-core processor.

  • Pros: It's incredibly easy. You don't have to change a single line of your application code.
  • Cons: There is a hard physical limit to how big a single machine can get. More dangerously, it is a Single Point of Failure (SPOF). If someone trips over the power cord of your supercomputer, your entire global application goes offline.

Horizontal Scaling (Scale Out)

Horizontal scaling means adding more computers. Instead of one massive server, you buy 100 cheap, standard servers and cluster them together.

  • Pros: Infinite scalability. If you need more power, just add more machines. High availability—if 3 servers catch fire, the other 97 keep the application running perfectly.
  • Cons: Immense complexity. Your code must now operate in a distributed environment. How do you route user traffic evenly? What if a user logs in on Server A, but their next request goes to Server B?

2. Load Balancing and Caching

To make Horizontal Scaling work, we need a traffic cop.

A Load Balancer sits in front of your cluster of web servers. When a user visits your site, they connect to the Load Balancer. The Load Balancer looks at the health of your servers and routes the request to the one with the least traffic. If a server dies, the Load Balancer instantly stops sending traffic to it.

Caching: The Ultimate Speed Boost

Querying a database is slow because it requires searching through disks. If a million users request the front page of your news site, generating that page via database queries a million times will crush your database.

Enter Caching (using in-memory stores like Redis or Memcached). RAM is orders of magnitude faster than a hard drive. When the first user requests the front page, the server asks the database, builds the page, and saves a copy in Redis. For the next 999,999 users, the server intercepts the request, grabs the pre-built page directly from Redis memory in a microsecond, and completely bypasses the database.


3. Database Scaling Strategies

Web servers are "stateless"—they don't store permanent data, making them easy to scale horizontally. Databases are "stateful"—they hold the truth. Scaling state is the hardest problem in computer science.

When a single database can no longer handle the load, we use two main strategies:

1. Read Replicas (Master-Slave Architecture)

In most applications (like Twitter or Wikipedia), there are vastly more "Reads" (viewing tweets) than "Writes" (posting a tweet). We can create multiple copies of the database. We designate one database as the Master and the others as Slaves (Read Replicas).

  • All Writes (Inserts, Updates, Deletes) must go to the Master.
  • The Master instantly broadcasts those changes to copy over to the Slaves.
  • All Reads (Selects) are routed by the application to the fleet of Slaves. This dramatically reduces the burden on the Master database.

2. Sharding (Partitioning)

What if the database is so massive (e.g., petabytes of data) that it physically cannot fit on the hard drive of a single Master server? We use Sharding. Sharding means cutting the database horizontally and putting different pieces on entirely separate servers.

  • Example: Server 1 holds users with last names A-M. Server 2 holds users with last names N-Z.
  • The Catch: Sharding adds nightmare-level complexity. If you want to run a query to "Find the oldest user on the platform", you have to query both databases and merge the results in your application code.

4. The Holy Grail: The CAP Theorem

In 2000, computer scientist Eric Brewer formulated the CAP Theorem, which states that in a distributed data system, it is mathematically impossible to guarantee more than two of the following three properties simultaneously:

  1. Consistency (C): Every read request receives the most recent, up-to-date write. (All nodes see the exact same data at the same time).
  2. Availability (A): Every request receives a non-error response. (The system is always online).
  3. Partition Tolerance (P): The system continues to operate even if the network fails and severs communication between nodes.

The Reality Check: In distributed systems (like cloud computing), network partitions (P) are a guarantee. Cables break, routers fail, and packets drop. Therefore, you must build for Partition Tolerance.

This means when a network failure occurs, as an architect, you are forced to choose between Consistency or Availability:

  • CP (Consistency over Availability): Example: A Banking System. If the network breaks between the ATM database and the central bank, the ATM will refuse to let you withdraw money (it sacrifices Availability). It prefers to be offline rather than risk giving you money you don't have (maintaining Consistency).
  • AP (Availability over Consistency): Example: Facebook Likes. If the network breaks, and you "like" a photo, Facebook will happily accept the like and show it to you (maintaining Availability). However, your friend in another state might look at the same photo and not see your like for a few minutes (sacrificing strict Consistency for "Eventual Consistency").

5. Microservices vs Monoliths

As a company grows, the way code is organized must change.

The Monolith: Everything (User Authentication, Payment Processing, Video Uploads) is packed into a single, massive codebase and deployed as one unit.

  • Pros: Simple to develop, easy to test, no network latency between functions.
  • Cons: When 100 developers work on one codebase, they step on each other's toes. If the Video Upload module has a memory leak and crashes, it takes the Payment Processing module down with it.

Microservices: The Monolith is shattered into dozens of tiny, independent applications. You have a "User Service", a "Payment Service", etc. Each is owned by a small team, deployed independently, and communicates with others via HTTP APIs or Message Queues.

  • Pros: Independent scaling. If the Payment Service is getting hammered, you can scale only the Payment Service without scaling the whole application. Fault isolation.
  • Cons: Monumental operational complexity. Debugging a request that jumps across 7 different microservices requires advanced tracing and observability tools.

6. Event-Driven Architecture

In a Microservice world, having Service A directly call Service B via an HTTP API creates tight coupling. If Service B goes offline, Service A starts failing too.

To fix this, we use Message Brokers (like Apache Kafka or RabbitMQ) to create an Event-Driven Architecture.

Instead of Service A calling Service B, Service A simply shouts into a megaphone (the Event Bus): "Hey, a new user registered!" Service A then goes back to its job. Service B (and Service C, and Service D) are subscribed to the Event Bus. When they hear the message, they process it at their own pace.

  • The Magic: If Service B is completely crashed, the message sits safely in the queue. When Service B reboots, it reads the queue and catches up. Zero data loss, and Service A never even knew Service B was down.

Conclusion

System Design is not about finding the "perfect" architecture; it's about choosing the right tradeoffs for your specific problem. An architecture that works for an e-commerce site will be disastrous for a real-time multiplayer game. By mastering load balancing, caching, database scaling, and decoupling via events, you gain the superpower to build software that can withstand the weight of the world.

Tags

#system-design#architecture#scalability