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_id or user_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:

  1. 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.
  2. Log Replication: The leader accepts write proposals, appends them to its log, and sends AppendEntries RPCs to follower nodes. Once a majority quorum acknowledges the log entry, the leader commits it to state machine.
  3. 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$:

Quorum Size Q = floor(N / 2) + 1 - In a 5-node cluster, Quorum Q = 3 nodes. - If network splits into 3-node and 2-node partitions: - The 3-node partition HAS quorum (3 >= 3) -> Remains operational. - The 2-node partition LACKS quorum (2 < 3) -> Rejects writes immediately.

6. Python Consistent Hashing Ring Implementation

Below is a Python implementation of a Consistent Hashing ring with virtual nodes (Vnodes) and MD5 hashing:

import hashlib import bisect from typing import List, Optional class ConsistentHashRing: def __init__(self, replica_vnodes: int = 100): self.replica_vnodes = replica_vnodes self.ring = [] # Sorted list of vnode hash integers self.ring_map = {} # Hash -> Physical Node string def _hash(self, key: str) -> int: md5_hex = hashlib.md5(key.encode('utf-8')).hexdigest() return int(md5_hex[:8], 16) # Convert first 8 hex chars to integer def add_node(self, node: str): for i in range(self.replica_vnodes): vnode_key = f"{node}-vnode-{i}" vnode_hash = self._hash(vnode_key) bisect.insort(self.ring, vnode_hash) self.ring_map[vnode_hash] = node def remove_node(self, node: str): for i in range(self.replica_vnodes): vnode_key = f"{node}-vnode-{i}" vnode_hash = self._hash(vnode_key) idx = bisect.bisect_left(self.ring, vnode_hash) if idx < len(self.ring) and self.ring[idx] == vnode_hash: del self.ring[idx] del self.ring_map[vnode_hash] def get_node(self, data_key: str) -> Optional[str]: if not self.ring: return None key_hash = self._hash(data_key) idx = bisect.bisect_right(self.ring, key_hash) if idx == len(self.ring): idx = 0 # Wrap around ring return self.ring_map[self.ring[idx]]

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.