6. Replication
Chapter 6 of Designing Data-Intensive Applications — 13 sections.
"The major difference between a thing that might go wrong and a thing that cannot possibly go wrong is that when a thing that cannot possibly go wrong goes wrong, it usually turns out to be impossible to get at or repair." — Douglas Adams
Replication = keeping a copy of the same data on multiple machines connected via a network.
Why:
- Latency — keep data geographically close to users
- Availability & durability — keep working even if parts have failed
- Read throughput — scale out the number of machines serving read queries
If the data doesn't change, replication is easy — copy it once and you're done. ALL the difficulty lies in handling CHANGES to replicated data.
Assumption for this chapter: the dataset is small enough that each machine holds a copy of the entire dataset. Ch 7 relaxes this (sharding).
Three families of algorithms — almost all distributed databases use one of these three:
SINGLE-LEADER MULTI-LEADER LEADERLESS
───────────── ──────────── ──────────
client ──▶ [L] client ──▶ [L₁] [L₂] ◀── client client ──▶ [R] [R] [R]
│ │ │ ╲ ╱ ◀── parallel
▼ ▼ ▼ ╲ ╱ writes AND reads
[F][F][F] ╳ to several nodes
╱ ╲
ONE node orders writes each leader also acts as NO leader; no ordering
followers apply in the a FOLLOWER to the others imposed; clients detect
SAME order and correct stale nodesThe principles haven't changed much since the 1970s, because the fundamental constraints of networks have remained the same. Nevertheless, concepts such as eventual consistency still cause confusion.
Sections
- 6.0Backups vs replication — they are NOT the same thing
- 6.16Single-Leader ReplicationThe leader logs every write statement it executes and sends the statement log to followers; each follower parses and executes that SQL as if received from a client.
- 6.21Problems with Replication LagAnd on where to solve it: dealing with these issues in application code is complex and easy to get wrong.
- 6.36Multi-Leader ReplicationBoth perform automatic merges for all the types above; different design philosophies and performance characteristics.
- 6.43Leaderless ReplicationIn some implementations the client sends writes to several replicas directly; in others a coordinator node does this on the client's behalf.
- 6.52Detecting Concurrent Writes
- 6.6Deep dives1Technology deep dives
- 6.7Failure catalogProduction failure catalog for this chapter
- 6.8Decision sheet1Decision cheat sheetAutomatic if your failover manager does real fencing and your timeout is tuned.
- 6.9Worked examplesWorked examplesThis is why quorums are seldom more than 4-of-7 or 5-of-9.
- 6.10Self-testSelf-test
- 6.11TerminologyTerminology introduced here
- 6.12Forward linksForward links