How Does Horizontal Scaling Work in Distributed Systems: A Complete Technical Guide

Horizontal scaling (scale-out) adds additional server instances to distribute workload rather than upgrading a single machine’s hardware, enabling distributed systems to handle millions of concurrent users through stateless services, load balancers, and data sharding.

In the liquidslr/system-design-notes repository, horizontal scaling is implemented across multiple architectural layers—from the web tier to the database layer. This guide examines the specific mechanisms, code patterns, and file structures that enable a service to grow from a single-machine prototype to a globally distributed architecture.

Stateless Services: The Foundation of Horizontal Scaling

Stateless architecture is a prerequisite for horizontal scaling. When session data resides in a shared datastore rather than local memory, any web server instance can handle any user request, allowing nodes to be added or removed without disrupting active sessions.

According to the Stateless Web Tier section in 01. Scaling/Readme.md, this pattern separates user state from application logic:

  • Session tokens are stored in Redis or similar distributed caches
  • Web servers authenticate requests against the shared store
  • Failed instances can be replaced instantly without data loss

This design ensures that adding a fourth server provides exactly four times the capacity, with no session affinity constraints tethering users to specific hardware.

Load Balancers: Distributing Traffic Across Instances

A load balancer sits between clients and the server pool, distributing incoming requests across available nodes. As documented in 01. Scaling/Readme.md, this component monitors server health and reroutes traffic when instances fail, preventing any single node from becoming a bottleneck.

The distribution algorithm determines how horizontal scaling translates to performance gains:

  • Round-robin: Cycles through servers sequentially
  • Least connections: Routes to the instance with fewest active requests
  • IP hash: Maintains session consistency for stateful legacy systems

When traffic spikes occur, new instances register with the load balancer, which immediately begins including them in the rotation without requiring configuration changes to client applications.

Database Sharding: Horizontal Scaling for Data Storage

While stateless web tiers scale easily, the data tier requires sharding (horizontal partitioning) to avoid becoming a chokepoint. The repository’s 01. Scaling/Readme.md explains that data is partitioned by a sharding key (e.g., user_id), with each shard hosted on a separate database server.

Key characteristics of this approach:

  • Each shard contains a subset of the total dataset
  • Queries without the sharding key must broadcast to all shards
  • Storage capacity grows linearly with the number of database nodes

This pattern allows the data layer to expand independently from the application layer, ensuring that storage I/O does not constrain overall system throughput.

Consistent Hashing: Minimizing Data Movement

When adding or removing database shards in a horizontally scaled system, consistent hashing minimizes the amount of data that must be relocated. The implementation details in 05. Consistent Hashing/Readme.md demonstrate how this algorithm maps data to a circular hash ring rather than using naive modulo arithmetic.

Benefits for distributed systems:

  • Adding a shard requires remapping only 1/n of the data (where n is the new shard count)
  • Removal of failed nodes redistributes load evenly across remaining instances
  • Virtual nodes prevent hotspots by distributing keys more uniformly

This technique is critical for maintaining performance during scaling events, preventing the "thundering herd" of cache misses that occurs when large portions of the dataset shift simultaneously.

Multi-Data-Center Deployment: Geographic Horizontal Scaling

True horizontal scaling extends beyond a single data center. The repository documents in 01. Scaling/Readme.md that Geo-DNS and data replication spread traffic across geographic regions, creating a global horizontal footprint.

This architecture provides:

  • Latency reduction by routing users to nearest data centers
  • Disaster recovery through cross-region replication
  • Compliance with data sovereignty requirements

Each data center operates as an independent scaling unit, allowing regional capacity to grow based on local demand patterns.

Practical Implementation Examples

The following Python implementations demonstrate the three core components of horizontal scaling as practiced in the liquidslr/system-design-notes repository.

Stateless Web Handler

This Flask application stores session data in Redis, ensuring any instance can authenticate any request:

from flask import Flask, request, jsonify
import redis

app = Flask(__name__)
store = redis.Redis(host='redis-cluster', port=6379)

@app.route('/login', methods=['POST'])
def login():
    user = request.json['user']
    # Authenticate logic here

    token = generate_token(user)
    # Store in shared Redis—accessible to all web servers

    store.set(token, user, ex=3600)
    return jsonify({'token': token})

Round-Robin Load Balancer

This simplified logic illustrates how a load balancer distributes requests:

servers = ['10.0.0.1', '10.0.0.2', '10.0.0.3']
counter = 0

def pick_server():
    global counter
    server = servers[counter % len(servers)]
    counter += 1
    return server

Sharding with Consistent Hashing

This function routes database writes to the appropriate shard:

import hashlib

def shard_for_key(key, num_shards):
    h = int(hashlib.sha256(str(key).encode()).hexdigest(), 16)
    return h % num_shards

def write_user_profile(user_id, data, db_shards):
    shard_id = shard_for_key(user_id, len(db_shards))
    db = db_shards[shard_id]
    db.put(user_id, data)

Summary

  • Horizontal scaling increases capacity by adding server instances rather than upgrading hardware, distributing load across multiple nodes in 01. Scaling/Readme.md.
  • Stateless services store session data externally (e.g., Redis), allowing any instance to handle any request and enabling seamless node replacement.
  • Load balancers distribute incoming traffic across the server pool, monitoring health and rerouting requests from failed instances as documented in the scaling guide.
  • Database sharding partitions data by key (such as user_id) across multiple database servers, preventing the storage tier from becoming a bottleneck.
  • Consistent hashing minimizes data movement when adding or removing shards, maintaining balanced load distribution with algorithms detailed in 05. Consistent Hashing/Readme.md.
  • Multi-data-center deployments extend horizontal scaling geographically through Geo-DNS and cross-region replication.

Frequently Asked Questions

What is the difference between horizontal and vertical scaling?

Horizontal scaling (scale-out) adds more machines to a pool of resources, while vertical scaling (scale-up) increases the CPU, RAM, or storage of a single existing server. According to the repository's 01. Scaling/Readme.md, horizontal scaling provides better fault tolerance and elasticity, whereas vertical scaling encounters hardware limits and creates single points of failure.

How does consistent hashing improve horizontal scaling?

Consistent hashing maps both data and cache servers to a hash ring, ensuring that when a node is added or removed, only a small fraction of keys need remapping. As implemented in 05. Consistent Hashing/Readme.md, this prevents the massive cache invalidation that occurs with traditional modulo-based sharding, maintaining system performance during scaling events.

Why must services be stateless for horizontal scaling?

Stateless services store no client-specific data in local memory between requests. This design, emphasized in the Stateless Web Tier section of 01. Scaling/Readme.md, allows any server instance to handle any request, enabling load balancers to distribute traffic evenly without session affinity constraints and allowing failed nodes to be replaced without data migration.

How does database sharding work with horizontal scaling?

Database sharding partitions data across multiple database instances based on a sharding key (e.g., user_id), with each shard residing on a separate server. This pattern, documented in 01. Scaling/Readme.md, allows the data tier to scale horizontally by adding database nodes independently of the application layer, though cross-shard queries require coordination across multiple instances.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →