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:
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.
“Each node may be the leader for some shards and a follower for other shards, but each shard still has only one leader.”
Terminology — the same thing under many names
| System | Name |
|---|---|
| Kafka | partition |
| CockroachDB | range |
| HBase, TiDB | region |
| Couchbase | vBucket |
| Riak | vnode |
| Cassandra | token-range |
| Bigtable, YugabyteDB, ScyllaDB | tablet |
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).
Sections
- 7.1Pros and Cons of Sharding
- 7.2Sharding for Multitenancy
- 7.39Sharding of Key-Value Data
- 7.42Request Routing
- 7.52Sharding and Secondary Indexes
- 7.6Deep divesTechnology deep dives
- 7.7Failure catalogProduction failure catalog for this chapter
- 7.8Decision sheet1Decision cheat sheetOnly if data volume or WRITE throughput exceeds one machine.
- 7.9Worked examplesWorked examples
- 7.10Self-testSelf-test
- 7.11TerminologyTerminology introduced here
- 7.12Forward linksForward links