Learn Labs
7. Sharding

7. Sharding

Chapter 7 of Designing Data-Intensive Applications — 12 sections.

"Clearly, we must break away from the sequential and not limit the computers. We must state definitions and provide for priorities and descriptions of data. We must state relationships, not procedures." — Grace Murray Hopper

A distributed database distributes data in two ways:

Replication (Ch 6) — store a copy of the same data on multiple nodes
copydatasetA B CnodeA B CnodeA B CnodeA B C
Sharding (Ch 7) — split the data into pieces, different pieces on different nodes
splitdatasetA B CnodeAnodeBnodeC
Figure 7.0.1A distributed database distributes data in two ways

Normally, shards are defined such that each piece of data (record, row, document) belongs to EXACTLY ONE shard. In effect, each shard is a small database of its own — although some systems support operations touching multiple shards at once.

Sharding is usually combined with replication, so copies of each shard live on multiple nodes: each record belongs to exactly one shard, but may still be stored on several nodes for fault tolerance.

NODE 1NODE 2NODE 3shard 1★ leadershard 1followershard 1followershard 2followershard 2★ leadershard 2followershard 3followershard 3followershard 3★ leader★ leader for that shard · every other copy is a follower

“Each node may be the leader for some shards and a follower for other shards, but each shard still has only one leader.”

Figure 7.0.2Sharding

Terminology — the same thing under many names

SystemName
Kafkapartition
CockroachDBrange
HBase, TiDBregion
CouchbasevBucket
Riakvnode
Cassandratoken-range
Bigtable, YugabyteDB, ScyllaDBtablet

PostgreSQL treats them as DISTINCT concepts: partitioning splits a large table into several files ON THE SAME MACHINE (advantages: e.g. very fast to delete an entire partition), whereas sharding splits a dataset ACROSS MULTIPLE MACHINES. In many other systems, partitioning is just another word for sharding.

Etymology, because it's a good story: one theory traces "shard" to the online RPG Ultima Online, in which a magic crystal was shattered into pieces, and each shard refracted a copy of the game world — so "shard" came to mean one of a set of parallel game servers, then carried over to databases. Another theory: an acronym for System for Highly Available Replicated Data, a 1980s database whose details are lost to history.

⚠️ Partitioning has NOTHING to do with NETWORK PARTITIONS (netsplits) — a type of fault in the network between nodes (Ch 9).


On this page