When a relational or NoSQL database grows beyond the CPU, RAM, or storage capacities of a single high-spec server, database architects must scale out horizontally. **Database Sharding** partitions data across multiple independent database nodes, while **Replication** copies data across nodes to guarantee high availability and fault tolerance.
However, scaling out storage introduces complex distributed systems challenges: rebalancing shards without downtime, maintaining data consistency during network partitions, avoiding split-brain scenarios, and resolving write conflicts.
1. Horizontal Sharding vs Vertical Partitioning
Understanding partition strategies is foundational to database scale:
- Vertical Partitioning: Splitting tables by columns (e.g., storing user login credentials in Table A on Node 1, and heavy user profile bio blobs in Table B on Node 2).
- Horizontal Sharding: Splitting table rows across independent database instances (Shards) based on a designated Shard Key (e.g.
tenant_idoruser_id). Each shard holds identical column schemas but different row subsets.
2. Consistent Hashing Ring Topology
A naive sharding formula uses modulo hashing: $\text{Shard} = \text{Hash}(K) \pmod N$, where $N$ is the number of database nodes.
The Flaw: If node count changes from $N=4$ to $N=5$ (adding a node to handle load), almost $100\%$ of keys remap to new shards! This triggers massive cluster-wide data re-shuffling that can cripple database performance.
Consistent Hashing solves this by mapping both database nodes and data keys onto a $360^\circ$ circular hash ring space ($0$ to $2^{32}-1$):
| Consistent Hashing Feature | Engineering Benefit |
|---|---|
| Ring Assignment | A key is assigned to the first database node encountered moving clockwise around the ring. |
| Node Addition / Removal | When adding or removing a node, only $K/N$ keys need to be migrated on average. Remaining nodes retain their existing keys. |
| Virtual Nodes (Vnodes) | Each physical node is assigned 100–250 virtual positions across the ring to ensure even data distribution and prevent hot-spotting. |
3. Database Replication Topologies
A. Active-Passive (Leader-Follower) Replication
All write queries execute on a single Leader node. The leader streams binary write logs (WAL / binlog) to Passive Follower nodes asynchronously or synchronously. Followers handle read queries exclusively.
B. Active-Active (Multi-Leader) Replication
Multiple master nodes accept writes simultaneously across different geographic data centers.
Conflict Resolution: Requires explicit strategies such as Last-Write-Wins (LWW based on NTP clocks), CRDTs (Conflict-Free Replicated Data Types), or application-level merge rules.
4. Distributed Consensus: The Raft Protocol
Distributed databases (such as CockroachDB, TiDB, and etcd) rely on the **Raft Consensus Protocol** to maintain a replicated state machine across independent nodes.
Raft decomposes consensus into three self-contained sub-problems:
- Leader Election: If follower nodes stop receiving heartbeats from the current leader within a randomized election timeout (150ms–300ms), they transition to Candidate state, increment term number, and request votes. The candidate receiving a majority quorum vote becomes the new Leader.
- Log Replication: The leader accepts write proposals, appends them to its log, and sends
AppendEntriesRPCs to follower nodes. Once a majority quorum acknowledges the log entry, the leader commits it to state machine. - Safety & Invariants: If a log entry is committed in a given term, that entry will be present in the logs of the leaders for all higher-numbered terms.
5. Resolving Split-Brain Scenarios
A **Split-Brain** occurs when a network partition cuts a cluster in half. If both halves assume the other half is dead and elect their own independent leaders, both accept writes simultaneously, corrupting data integrity.
Mitigation: Strict Majority Quorum Math
For a cluster of $N$ nodes, a leader can only commit writes or remain operational if it communicates with a strict majority quorum $Q$:
6. Python Consistent Hashing Ring Implementation
Below is a Python implementation of a Consistent Hashing ring with virtual nodes (Vnodes) and MD5 hashing:
7. Summary Guidelines
- Choose High-Cardinality Shard Keys: Select shard keys with high cardinality (like
uuid) to prevent uneven data skew. - Always Use Odd Node Counts: Deploy 3, 5, or 7 nodes for Raft/Paxos consensus clusters so split-brain network partitions yield a single clear majority partition.
Join the Technical Discussion
Have questions about this architecture? Drop a comment below.