Horizontal Scaling (Scaling Out)

⭐ Interview Importance: HIGH
⏱️ Revision Time: 3 min

Concept

When you hit the absolute physical ceiling of Vertical Scaling (you cannot buy a bigger server), or the price becomes astronomically high, you must switch to Horizontal Scaling (Scaling Out).

Horizontal Scaling means adding more servers, rather than bigger servers.
Instead of one 10,000supercomputer,youbuyten10,000 supercomputer, you buy ten 1,000 standard computers, connect them via a high-speed network, and distribute the database load across all of them.

The Node.js Advantage vs The Database Problem

Horizontal scaling is incredibly easy for Backend Web Servers (Node.js). Because Node.js APIs are Stateless (they don’t save data on their own hard drives), you can just spin up 50 identical Docker containers behind an AWS Load Balancer. If Container #4 crashes, the Load Balancer just routes traffic to Container #5. Zero data is lost.

Horizontal scaling is notoriously difficult for Databases because databases are Stateful.
If you just spin up 5 identical PostgreSQL servers, and User A sends an INSERT command to Database Server #1, the other 4 Database Servers have no idea that the user was created. They are completely out of sync.

To Horizontally Scale a relational database, you must implement complex data synchronization strategies. The two primary strategies are:

  1. Read Replicas (Primary/Replica Architecture)
  2. Database Sharding

The Trade-Offs of Horizontal Scaling

  • Pros:
    • Infinite Scalability: You can theoretically add 10,000 servers to a cluster, handling trillions of requests per second (like Google Spanner or Cassandra).
    • Fault Tolerance (High Availability): If one server’s motherboard burns out, the other 9 servers in the cluster instantly take over the traffic. The application never goes offline.
  • Cons:
    • Architectural Complexity: Your application must now be aware of the network topology. Which server do I read from? Which server do I write to?
    • Eventual Consistency: When data is duplicated across 10 servers, there is a physical speed-of-light delay in synchronizing them. A user might write data to Server 1, instantly read from Server 2, and the data won’t be there yet.

Interview Questions

Q: In cloud environments (like AWS Auto-Scaling Groups), your Node.js servers can automatically scale horizontally based on CPU usage. If traffic spikes, AWS boots up 10 new Node.js servers automatically. Can you do the exact same thing with a PostgreSQL database?
A: No. You cannot easily “auto-scale” a standard relational database horizontally in real-time.
Booting up a new Node.js server takes 2 seconds because it is stateless. Booting up a new PostgreSQL server requires physically copying the entire 500GB database file from the primary server to the new server before it can safely accept traffic. This synchronization can take hours. Relational databases must be scaled proactively, not reactively. (Note: Specialized Cloud-Native databases like Amazon Aurora and Serverless SQL have heavily mitigated this by separating the compute tier from the storage tier).