Command Palette

Search for a command to run...

Hectal
Phase 9Advanced10 of 18 in Database Design

Replication, Partitioning and Sharding

Primary/replica replication, lag and read-your-writes, failover and split brain; range, list, hash and time partitioning with pruning and retention; shard keys, resharding, cross-shard queries and hotspot mitigation.

Scale in the cheapest order: tune queries and indexes, add cache and read replicas, partition big tables, and shard only when a single primary can't handle the writes or the data. Each step adds operational cost and new failure modes.

0/5 · 0%
5 topics ~43 min 7 code blocks & diagrams
Start with the first topic
1
9.1

Replication: Replicas, Lag, Consistency and Failover

A primary streams its WAL to replicas that replay it. Replicas scale reads and provide failover targets. Asynchronous replication is fast but can lose recent commits on failover and serves stale reads; synchronous replication waits for replica confirmation. Failover needs fencing to avoid split brain, and applications need read-your-writes routing.

9 min 1 diagram 1 code practice

2
9.2

Partitioning: Range, List, Hash and Time

Partitioning splits one large table into smaller physical tables on one server by a partition key. Queries that filter on the key touch only relevant partitions (pruning), old data is dropped by detaching a partition instead of deleting rows, and maintenance runs per partition. It helps manageability far more than raw speed, and a bad key makes everything slower.

9 min 1 code practice

3
9.3

Sharding: Shard Keys and Strategies

Sharding spreads data across multiple independent database servers so writes and storage scale horizontally. The shard key decides everything: it should spread load evenly, keep each request's data on one shard, and never need to change. Strategies are hash, range, directory (lookup) and geographic, usually with many logical shards mapped onto fewer physical nodes.

9 min 2 code practice

4
9.4

Resharding, Cross-Shard Queries and Global Indexes

Once sharded, the hard parts are queries that don't include the shard key (scatter-gather or a global secondary index), transactions spanning shards (avoid, or use sagas or 2PC), and resharding live data (copy, catch up via change streams, switch the routing mapping, then clean up).

8 min 1 diagram practice

5
9.5

Hotspot Analysis and Mitigation

Hotspots are rows, keys, partitions, shards, indexes or nodes that receive a disproportionate share of traffic: a celebrity account, a flash-sale product, today's partition, the right edge of a sequential index. Find them with per-key metrics, then mitigate with caching, key splitting (salting), write buffering, replication, or isolating the hot entity.

8 min 1 code practice