What is Sharding, Partitioning and Why Sharding Managers?

💡 Originally published on Medium.
TL;DR
When database workloads outgrow the vertical CPU, memory, and disk capacity of a single machine, systems must partition data horizontally across multiple database instances—a pattern known as Sharding. However, managing sharded architectures introduces severe operational complexities: cross-shard joins, schema migrations, shard rebalancing, and routing logic. Sharding Managers (e.g., Vitess, Citus, Spanner) abstract these distributed challenges away from application code.
Horizontal Sharding vs. Vertical Partitioning
graph TD
subgraph Monolith["Single Large Database (Scaling Bottleneck)"]
DB["Monolithic Database<br/>(Users, Orders, Inventory, Payments)"]
end
subgraph VertPart["Vertical Partitioning (Split by Domain/Columns)"]
V1["Users DB"]
V2["Orders DB"]
V3["Payments DB"]
end
subgraph HorizShard["Horizontal Sharding (Split Rows by Shard Key)"]
H1["Shard 1 (User IDs 1 - 1,000,000)"]
H2["Shard 2 (User IDs 1,000,001 - 2,000,000)"]
H3["Shard 3 (User IDs 2,000,001 - 3,000,000)"]
end
DB -->|Split by Domain| VertPart
DB -->|Split Rows by Hash/Range| HorizShard
- Vertical Partitioning: Splitting a schema by tables or columns into distinct microservice databases (e.g., separating
User ProfilesfromAudit Logs). - Horizontal Partitioning (Sharding): Keeping identical table schemas across multiple database servers, where each server holds a subset of total rows based on a Shard Key (e.g.,
user_idortenant_id).
Sharding Strategies
1. Hash-Based (Consistent Hashing)
A hash function maps the shard key to a discrete shard index:
$$ ext{Shard Index} = ext{Hash}( ext{Shard Key}) \pmod N$$
- Pros: Uniform data distribution; avoids hotspot creation.
- Cons: Adding new shard nodes requires complex data migration and resharding.
2. Range-Based Sharding
Data is partitioned according to contiguous key ranges (e.g., A–F, G–M, N–Z, or timestamps):
- Pros: Simple to query and range-scan.
- Cons: High risk of write hotspots (e.g., recent date partitions receiving all active writes).
Why You Need a Sharding Manager
Implementing sharding logic directly inside application code tightly couples business services to database topology. A Sharding Manager acts as an intelligent proxy layer:
graph LR
App["Application Service"] --> SM["Sharding Manager / Router Proxy<br/>(Query Parsing & Routing Engine)"]
SM -->|Hash(user_id) % 3 == 0| S1["Database Shard 1"]
SM -->|Hash(user_id) % 3 == 1| S2["Database Shard 2"]
SM -->|Hash(user_id) % 3 == 2| S3["Database Shard 3"]
Core Responsibilities of Sharding Managers:
- Transparent Query Routing: Applications issue standard SQL queries; the manager routes reads and writes to the correct shard.
- Scatter-Gather Execution: Aggregates and merges multi-shard queries (e.g.,
SELECT COUNT(*) FROM orders) concurrently. - Dynamic Online Resharding: Splits and migrates shards in the background without incurring application downtime.
- Health & Connection Pooling: Manages thousands of pooled client connections to prevent database connection starvation.