Skip to main content

Database Sharding & Partitioning

Who this guide is for

What is Partitioning?

Partitioning splits a large dataset into smaller, more manageable pieces. There are two fundamentally different ways to split:

Vertical partitioning

Split a table by columns. Move rarely-accessed or large fields into a separate table while keeping hot fields together.

Vertical Partitioning Mechanics
Single Wide Table (Before Split)
user_idusernameemailprofile_bioavatar_blob
1user_1usr1@mail"Senior Java Dev..."0xFD29B8C...
2user_2usr2@mail"Senior Java Dev..."0xFD29B8C...
3user_3usr3@mail"Senior Java Dev..."0xFD29B8C...
โš ๏ธ Performance Limitation (Wide Table)
  • **Row Size Bloat**: Reading rows loads the massive `profile_bio` text and `avatar_blob` binary data into memory.
  • **Cache Pollution**: Fewer rows fit in the database buffer pool page, increasing disk I/O reads.
  • **Memory consumption**: 1,000 login query records pull ~25MB of unused binary chunks into memory page caches.

When to use it: when a small number of columns are accessed in 95% of queries and the rest inflate the row size, bloating cache and increasing I/O.

Horizontal partitioning (sharding)

Split a table by rows. Different subsets of rows live on entirely different database servers.

Horizontal Partitioning (Row-Based Sharding Modulo)
Shard RouterHash: 120 % 3 = 0Shard A (user_id % 3 = 0)๐Ÿ“ Storing user_id 120Shard B (user_id % 3 = 1)listeningShard C (user_id % 3 = 2)listening
๐Ÿ’ก **Modulo Sharding Mechanics:** Each row is routed entirely to a dedicated physical server partition based on `user_id % N`. While simple and effective, a major drawback is that changing `N` (adding/removing shards) forces almost all key locations to remap, causing massive data reshuffling.

When we say "sharding" in system design, we almost always mean horizontal partitioning. The rest of this guide focuses on it.


Why Shard?

A single database server has hard physical limits. Understanding which limit you're hitting determines the right solution:

BottleneckSymptomFirst-line solutionWhen you need sharding
Read loadSlow SELECTs, high CPU on readsAdd read replicasWhen even read replicas can't absorb the fan-out
Write throughputReplication lag, high disk I/OTune indexes, batch writesWhen write IOPS saturate the primary's disk
Dataset sizeDisk full, slow full-table scansArchival, compressionWhen data exceeds one server's storage capacity
RAM (working set)High disk reads, cache miss rate risingBigger instanceWhen the working set no longer fits in memory
Exhaust simpler options first

Before sharding, try in this order: query tuning โ†’ indexes โ†’ vertical scaling โ†’ read replicas โ†’ caching (Redis). Sharding adds enormous operational complexity. It should be a last resort, not a first instinct.

Concrete thresholds to consider sharding

  • Table rows exceeding 500Mโ€“1B and growing fast
  • Write throughput exceeding ~10K writes/sec sustained on a single primary
  • Dataset exceeding 2โ€“4 TB (point where even SSDs become expensive and slow)
  • P99 write latency increasing despite hardware upgrades

Sharding Strategies

Choosing a routing algorithm is the most important decision in your sharding design. Each strategy makes a different trade-off between write distribution, range query efficiency, and rebalancing cost.

1. Range partitioning

Divide data into continuous ranges of the shard key. Each shard owns a contiguous slice.

Shard 1: user_id 1 โ€“ 10,000,000
Shard 2: user_id 10,000,001 โ€“ 20,000,000
Shard 3: user_id 20,000,001 โ€“ 30,000,000
-- PostgreSQL declarative range partition
CREATE TABLE orders (
id BIGSERIAL,
user_id BIGINT NOT NULL,
created_at TIMESTAMPTZ NOT NULL,
total NUMERIC(12,2)
) PARTITION BY RANGE (created_at);

CREATE TABLE orders_2024_q1 PARTITION OF orders
FOR VALUES FROM ('2024-01-01') TO ('2024-04-01');

CREATE TABLE orders_2024_q2 PARTITION OF orders
FOR VALUES FROM ('2024-04-01') TO ('2024-07-01');

Pros:

  • Range queries are extremely efficient โ€” WHERE user_id BETWEEN 100 AND 500 touches exactly one shard.
  • Contiguous scans for reporting and archival are fast.
  • Easy to reason about which shard holds which data.

Cons:

  • High risk of hot spots (see below). Sharding by created_at or auto-increment means all new writes always go to the last shard.
  • Uneven data distribution if the key is not uniformly spread.
The hot spot problem with range partitioning

If you shard by timestamp or a monotonically increasing ID:

t=0: Shard 1 โ† all writes (idle: Shard 2, 3)
t=1: Shard 2 โ† all writes (idle: Shard 1, 3)
t=2: Shard 3 โ† all writes (idle: Shard 1, 2)

You've bought yourself horizontal hardware but not horizontal write throughput. The active shard is always the bottleneck. Mitigations: add a random prefix to the key, shard by a different key (e.g. user_id instead of order_time), or use consistent hashing.


2. Hash partitioning

Apply a hash function to the shard key and use modulo to assign a shard:

shard_id = hash(user_id) % num_shards

hash("user_1001") % 3 = 0 โ†’ Shard A
hash("user_1002") % 3 = 2 โ†’ Shard C
hash("user_1003") % 3 = 1 โ†’ Shard B

Pros:

  • Even data distribution โ€” hash functions spread keys uniformly, eliminating hot spots.
  • Simple to implement and reason about.

Cons:

  • Range queries become scatter-gather operations. user_id BETWEEN 1 AND 1000 hits all shards because adjacent keys hash to different shards.
  • Rebalancing is catastrophic. Adding a shard changes num_shards from N to N+1, so hash(key) % (N+1) gives a different result for almost every key. Nearly 100% of data must migrate.
3 shards โ†’ 4 shards:
key "abc": hash % 3 = 1 (was Shard B) โ†’ hash % 4 = 3 (now Shard D)
key "xyz": hash % 3 = 0 (was Shard A) โ†’ hash % 4 = 1 (now Shard B)
... virtually every key must move

This is the fundamental weakness that consistent hashing was designed to solve.


3. Consistent hashing

Rather than modulo arithmetic, both keys and nodes are mapped onto a virtual ring from 0 to 2ยณยฒโˆ’1. A key is owned by the first node clockwise from its hash position.

Basic Consistent Hashing Ring
0%50%N1 (25%)N3 (75%)10%40%60%
โš™ Key-to-Node Routing
๐Ÿ’ก Click a key percentage on the ring or choose from the list to trace its routing clockwise to the owning node.

Adding a node: only the keys between the new node and its predecessor on the ring need to move. ~1/N keys migrate instead of ~100%.

Consistent Hashing Ring Topology Updates
N1 (25%)N2 (50%)N3 (75%)N4 (60%)A (40%)B (55%)C (70%)
๐Ÿ“ˆ Migration Assessment
Key A (40%): Assigned to N2. Unchanged โœ…
Key B (55%): Assigned to N3 (75%).
Key C (70%): Assigned to N3. Unchanged โœ…
Migration Footprint:
Click "Add Node N4" to trace key re-allocation.

Removing a node (failure or decommission): only that node's keys move to its successor. All other nodes are unaffected.

Pros:

  • Minimal data movement on topology changes (~1/N keys migrate).
  • Naturally supports heterogeneous nodes via virtual nodes.
  • No single point of failure โ€” ring is fully decentralised.

Cons:

  • Basic consistent hashing can still produce uneven distribution if nodes land unluckily on the ring.
  • More complex to implement than modulo hashing.
  • Non-uniform load is solved with virtual nodes (see senior section below).

4. Directory / lookup service

A central routing table (sometimes called a shard map) explicitly records where each key or key range lives:

Directory / Lookup Service Router Flow
ClientRouterShard Map DBShard 1Shard 21. Query key "usr_101"2. Directory Lookup3. Return Host Info4. Direct Routing5. Return Data
Click an arrow or click "Animate Flow" to trace step lookups in a directory service routing topology.

Pros:

  • Maximum flexibility โ€” you can manually move individual hot tenants to dedicated hardware.
  • Supports complex routing logic (geo-based, tier-based).
  • Shard rebalancing requires only updating the directory, not rehashing.

Cons:

  • Central directory is a single point of failure (mitigate with replication + caching).
  • Adds a network hop to every query.
  • Directory can become stale โ€” cache invalidation is a hard problem.

Strategy comparison

RangeHash moduloConsistent hashDirectory
Range queriesโœ… Efficient (single shard)โŒ Scatter-gatherโŒ Scatter-gatherโœ… If mapped
Write distributionโŒ Hot spots (monotonic keys)โœ… Evenโœ… Evenโœ… Flexible
Rebalancing cost๐ŸŸก Migrate rangeโŒ Migrate ~100%โœ… Migrate ~1/Nโœ… Update table
ComplexityLowLowMediumHigh
Used byPostgreSQL, MySQL, CassandraRedis Cluster (basic)Cassandra, DynamoDB, RiakVitess, some custom systems

Consistent Hashing โ€” Deep Dive

Virtual nodes (vnodes)

Basic consistent hashing with one point per physical node produces uneven distribution โ€” nodes land at random positions and some get much larger arcs than others.

The solution: each physical node claims multiple positions on the ring (virtual nodes). A cluster with 3 servers and 150 vnodes per server has 450 points on the ring, producing near-uniform distribution.

Consistent Hashing: Virtual Nodes (Vnodes)
S1_1S2_1S3_1S1_2S2_2S3_2S1_3S2_3S3_3S1_4S2_4S3_4
๐ŸŸข Uniform Load (Stable)
  • **Balanced Arcs**: Interleaving 12 vnode positions splits the ring into tiny slices, smoothing out load deviation.
  • **Proportional Capacity**: Heterogeneous servers can be scaled easily (e.g. S1 gets 256 vnodes, S3 gets 128 vnodes).
  • **Fault Isolation**: When S1 fails, its workload is distributed among S2 and S3 vnode slots on the ring instead of crashing a single neighbour.

Benefits of vnodes:

  • Even load distribution regardless of how many nodes join or leave.
  • A more powerful server can be given more vnodes to handle a proportionally larger share.
  • When a node fails, its vnodes are spread across many other nodes, distributing the recovery load instead of dumping everything on one neighbour.

Used by: Cassandra (default 256 vnodes/node), DynamoDB internally, Riak.

Replication with consistent hashing

In production, keys are not stored on just one node. They are replicated to the next N nodes clockwise on the ring (the replication factor):

Consistent Hashing Replication (Replication Factor = 3)
Node A (Primary)Node B (Replica 1)Node C (Replica 2)key "usr_99"
๐Ÿ”’ Durability Audit (RF=3)
Primary Assignment: Key `"usr_99"` hashes to 45ยฐ, traveling clockwise to first node **Node A**. Node A functions as coordinator.
Clockwise Replication: The write is synchronously/asynchronously streamed from Node A โ†’ **Node B** โ†’ **Node C**.
Failure Resilience: If Node A crashes, Node B takes over read/write queries. The data survives the concurrent loss of up to 2 replica servers.

This means every piece of data survives the loss of up to N-1 nodes in its replica group.


Cross-Shard Problems

Sharding is not free. It introduces a category of distributed systems problems that don't exist on a single server.

Cross-shard JOINs

SQL JOINs assume data lives in the same process. When tables are split across servers, you can't do a SQL JOIN across shards.

-- Works on a single DB:
SELECT u.name, o.total
FROM users u JOIN orders o ON u.id = o.user_id
WHERE u.country = 'VN';

-- With sharding: users and orders may be on different shards.
-- This JOIN is impossible at the DB layer.

Solutions:

Design your shard key so related entities always land on the same shard. If you shard users and orders both by user_id, all of Alice's data is on the same shard:

Shard A: users where user_id % 3 = 0
orders where user_id % 3 = 0
โ†’ JOIN between Alice's user row and Alice's orders never crosses a shard

Global transactions (distributed ACID)

A transaction that touches data on multiple shards cannot use standard ACID guarantees without a distributed coordination protocol.

Transfer $100 from Alice (Shard A) to Bob (Shard B):

Step 1: Debit Alice on Shard A โœ…
Step 2: Credit Bob on Shard B โŒ (Shard B crashes)

Result: Alice lost $100, Bob received nothing โ†’ inconsistency

Solutions:

ApproachHow it worksTrade-offs
Two-Phase Commit (2PC)Coordinator asks all shards to "prepare", then commits if all agree.Synchronous, slow, coordinator is a SPOF. Blocks if coordinator crashes mid-transaction.
Saga patternBreak transaction into a sequence of local transactions. Each step publishes an event. Compensating transactions undo steps on failure.Eventual consistency, complex to implement, no atomicity guarantee.
Avoid cross-shard writesDesign data model so all writes for a logical operation touch one shard only (co-location).Best option when achievable โ€” zero coordination cost.
Optimistic locking + reconciliationAllow eventual inconsistency; detect and reconcile conflicts later.Works for some use cases (analytics); unacceptable for money.
The best solution is avoiding the problem

If you find yourself writing a lot of cross-shard transactions, your shard key is probably wrong. Revisit whether a different key can co-locate the data that is written together.


Unique ID generation

Auto-incrementing primary keys (AUTO_INCREMENT, SERIAL) don't work across isolated shards โ€” each shard would generate its own id=1, id=2, creating duplicates when merged.

Solutions:

import uuid
id = uuid.uuid4() # e.g. "550e8400-e29b-41d4-a716-446655440000"

UUID v4: random 128-bit โ€” globally unique, but not sortable. Causes random index insertions, defeating B-tree locality.

UUID v7: timestamp-prefixed โ€” globally unique AND sortable by creation time. Preferred over v4 for database IDs.


Scatter-gather queries (no shard key in filter)

When a query doesn't include the shard key in its WHERE clause, the router has no choice but to send it to all shards and merge the results:

-- Shard key is user_id. This query has no user_id:
SELECT * FROM users WHERE email = '[email protected]';
-- โ†’ sent to ALL shards, results merged, sorted, returned

This is called a scatter-gather or fan-out query. At 10 shards it's 10x overhead. At 100 shards it's 100x.

Solutions:

SolutionDescriptionTrade-off
Global secondary indexMaintain a separate index table mapping non-key attributes to shard keysExtra write on every insert/update; index must be highly available
Mapping table in Redisemail โ†’ user_id stored in Redis; look up user_id first, then query correct shardExtra network hop; cache invalidation
Dual-write to a search indexWrite to both your DB and Elasticsearch on every mutationEventual consistency; more complex writes
Covering shard keyChoose a shard key that appears in all hot query patterns (e.g. tenant_id in a SaaS app)May not be possible for all queries

Rebalancing Strategies

When your cluster grows (adding nodes) or shrinks (node failure, decommission), data must be redistributed. How you do this determines downtime and migration cost.

Manual rebalancing

An operator explicitly decides which partitions to move and executes the migration. The system makes no automatic decisions.

  • Pro: full operator control, no surprise data movements.
  • Con: tedious and error-prone at scale; requires careful sequencing.

Used by: MongoDB (manual chunk migration via moveChunk).

Automatic rebalancing

The system continuously monitors shard load and migrates partitions in the background when imbalance is detected.

Cassandra:
โ†’ monitors each vnode's data volume
โ†’ auto-streams data to new nodes when they join
โ†’ background streaming, no downtime

DynamoDB:
โ†’ fully managed; partitions split/merge transparently
โ†’ no operator involvement required
  • Pro: hands-off operation; reacts to organic growth.
  • Con: background migrations consume I/O; can degrade query performance during rebalancing.

Fixed partition count

Instead of changing the number of partitions when nodes change, use a large fixed number of partitions (e.g. 1000) and redistribute ownership of those partitions to nodes:

1000 fixed partitions, 3 nodes:
Node A โ†’ partitions 0โ€“332
Node B โ†’ partitions 333โ€“666
Node C โ†’ partitions 667โ€“999

Adding Node D:
Node A โ†’ partitions 0โ€“249 (moved 83 partitions to D)
Node B โ†’ partitions 333โ€“582 (moved 84 partitions to D)
Node C โ†’ partitions 667โ€“916 (moved 83 partitions to D)
Node D โ†’ partitions 250โ€“332, 583โ€“666, 917โ€“999

Only metadata (partitionโ†’node mapping) changes; the partitions themselves remain stable. Used by Elasticsearch (fixed primary shards).

You can't change the number of primary shards in Elasticsearch after index creation

This is why choosing the right initial shard count matters so much. Rule of thumb: aim for shards of 10โ€“50 GB each. For a 500 GB index, 10โ€“50 shards is reasonable.


Scatter-Gather vs. Co-located Queries

Understanding when queries cross shards is essential for performance planning:

Query Partitioning: Co-located vs. Scatter-Gather
ClientQuery Routerusr_id filter presentShard AShard B (Target)Shard C๐ŸŸข 1 Network Hop
โšก Co-located Query (Single-Shard Routing)
  • **Shard Key Filter**: The SQL query contains `WHERE user_id = 101`. The Router hashes this key and routes to Shard B directly.
  • **Latency Profile**: Direct 1-network RTT routing. No cross-shard aggregation or sorting needed.
  • **Scalability**: Highly scalable. Doubling shards has zero impact on query execution latency.

Real latency impact: at P50 scatter-gather is ~3x slower; at P99 it's often 10โ€“20x slower because you wait for the slowest shard (the "long tail" problem).


Real-World Database Comparison

Strategy: Consistent hashing with vnodes (default 256 per node).

Shard key: The partition key in the PRIMARY KEY definition.

CREATE TABLE orders_by_user (
user_id UUID,
created_at TIMESTAMP,
order_id UUID,
total DECIMAL,
PRIMARY KEY ((user_id), created_at, order_id)
-- ^^^^^^^^^
-- partition key โ†’ determines shard
-- created_at, order_id โ†’ clustering keys (sort within shard)
) WITH CLUSTERING ORDER BY (created_at DESC);

Rebalancing: Automatic streaming when nodes join or leave. New nodes bootstrap data from neighbours.

Cross-shard JOINs: Not supported. You must denormalise into query-specific tables.


Sharding in SQL Databases โ€” Patterns

PostgreSQL: list partitioning

Useful when the partition key is categorical (e.g. region, status):

CREATE TABLE orders (
id BIGSERIAL,
region TEXT NOT NULL,
total NUMERIC(12,2)
) PARTITION BY LIST (region);

CREATE TABLE orders_apac PARTITION OF orders
FOR VALUES IN ('VN', 'SG', 'TH', 'PH', 'ID');

CREATE TABLE orders_emea PARTITION OF orders
FOR VALUES IN ('DE', 'FR', 'GB', 'NL');

CREATE TABLE orders_amer PARTITION OF orders
FOR VALUES IN ('US', 'CA', 'BR', 'MX');

Partition pruning

When the WHERE clause includes the partition key, PostgreSQL eliminates irrelevant partitions from the query plan โ€” only scanning the partitions that could contain matching rows:

EXPLAIN SELECT * FROM orders WHERE region = 'VN' AND total > 1000;
-- โ†’ Seq Scan on orders_apac (partition pruning eliminates orders_emea, orders_amer)

Sub-partitioning (composite)

Partitions can themselves be partitioned โ€” range by year, then hash within the year:

CREATE TABLE events (
id BIGSERIAL,
created_at DATE NOT NULL,
user_id BIGINT NOT NULL
) PARTITION BY RANGE (created_at);

CREATE TABLE events_2024 PARTITION OF events
FOR VALUES FROM ('2024-01-01') TO ('2025-01-01')
PARTITION BY HASH (user_id); -- sub-partition by user

CREATE TABLE events_2024_p0 PARTITION OF events_2024
FOR VALUES WITH (MODULUS 4, REMAINDER 0);
-- ... p1, p2, p3

Monitoring a Sharded Cluster

A sharded system introduces failure modes that don't exist on single-server databases. These are the key signals to watch:

MetricWhat it revealsAlert threshold (example)
Shard data size varianceUneven distribution โ€” one shard getting all writes>20% deviation from mean
Replication lag per shardA replica is falling behind; stale reads possible>30 seconds
Scatter-gather query ratio% of queries without shard key โ€” cross-shard overhead>5% of query volume
P99 cross-shard query latencyLong-tail queries from multi-shard fan-out>1 second
Chunk migration rate (MongoDB)Background rebalancing consuming I/OMonitor spikes during business hours
Hot partition rate (DynamoDB)One partition key receiving disproportionate trafficThrottled requests per partition

Interview Questions

For new learners

Q: What is the difference between replication and sharding?

Replication copies the same data to multiple nodes to improve read scalability and high availability (fault tolerance). Sharding splits different data across multiple nodes to improve write scalability and handle datasets too large for one disk. They are almost always used together: a distributed database will have multiple shards, and each shard will have multiple replicas.

Q: What is a hot spot in sharding, and how do you prevent it?

A hot spot occurs when one shard receives disproportionately more traffic than others. The most common cause is choosing a monotonically increasing shard key (timestamp, auto-increment ID) โ€” all new writes always go to the "latest" shard. Mitigations: choose a high-cardinality, randomly distributed shard key (like user_id), use hash partitioning, or add a random salt prefix to the key at the cost of scatter-gather on reads.

Q: When should you NOT shard?

When simpler alternatives can solve the problem: query tuning, adding indexes, vertical scaling (bigger server), adding read replicas, or in-memory caching. Sharding adds massive operational complexity โ€” cross-shard JOINs, distributed transactions, complex deployments, harder debugging. Only shard when you've genuinely exhausted the alternatives and benchmarks confirm the bottleneck.


For senior engineers

Q: Why is hash modulo (hash(key) % N) a bad strategy for dynamically scaling a database?

When you add one shard (N becomes N+1), the modulo result changes for nearly every key โ€” hash(key) % N and hash(key) % (N+1) rarely agree. This means close to 100% of data must migrate across the network to restore consistency. For a 10 TB dataset, this is a near-total reshuffle with massive I/O and potential downtime. Consistent hashing solves this by only requiring ~1/N of data to migrate when adding one node.

Q: How do virtual nodes (vnodes) improve consistent hashing?

Basic consistent hashing places one point per physical node on the ring. With three nodes, one might get a 40% arc, another 35%, another 25% โ€” uneven load. Virtual nodes assign each physical node 100โ€“300 positions on the ring, producing near-uniform distribution. Additionally: a stronger server can be given more vnodes to carry a larger share; when a node fails, its vnodes scatter to many different neighbours, distributing recovery load instead of overloading one node.

Q: You shard users by user_id. A user logs in with only their email. How do you find the right shard?

email is not the shard key, so the router doesn't know which shard to query. Two approaches: (1) Scatter-gather โ€” send the query to all shards and merge results. Works but is O(shards) in cost and unacceptable at scale. (2) Mapping table โ€” maintain a high-availability lookup (e.g. Redis or a small dedicated DB) that maps email โ†’ user_id. Every login first queries the mapping table to get user_id, then routes to the correct shard. The mapping table is tiny (just email + user_id) and can be replicated globally for low latency.

Q: How does the Saga pattern replace distributed transactions in a sharded system?

A Saga breaks a cross-shard operation into a sequence of local transactions, each publishing a domain event on completion. Downstream services listen and execute their step. If any step fails, compensating transactions run in reverse order to undo completed steps (e.g. refund a payment if inventory reservation fails). The key difference from 2PC: Saga is asynchronous and eventually consistent โ€” there is a window where the system is partially updated. This is acceptable for many business workflows (order fulfilment) but not for hard financial consistency (bank transfers).

Q: A DynamoDB table's user_id partition is throttled because one user generates 10x the normal traffic. How do you fix it?

The hot partition problem. Options: (1) Write sharding โ€” append a random suffix (user_id#0 through user_id#7) to distribute writes across 8 partitions, then scatter-gather on reads and aggregate. (2) Caching โ€” put a Redis cache in front for read-heavy hot users. (3) DAX (DynamoDB Accelerator) โ€” AWS in-memory cache for DynamoDB, microsecond reads for hot keys. (4) Data model redesign โ€” if this user's data is always accessed differently from others, consider a separate table or a different access pattern that doesn't funnel through one partition.

Q: Walk through how Cassandra rebalances data when a new node joins.

When a new node joins, it announces itself to the ring with a chosen token (or multiple tokens if using vnodes). Cassandra's gossip protocol propagates its existence to the cluster. The node then streams its responsible token ranges from its neighbours โ€” specifically, it takes ownership of its vnodes' key ranges, and the previous owners stream the corresponding data to it. Read and write traffic gradually shifts to the new node as the ring topology updates. The process is online โ€” the cluster continues serving traffic throughout. Once streaming completes, the operator can nodetool cleanup on the previous owners to reclaim disk space from keys that are no longer their responsibility.


See Also

๐Ÿ“–
Track Page Progress0 / 635 Read
Knowledge Base Completion0%