Lesson 23 / 25
Replication, Partitioning and Distributed Transactions
Describe how databases scale across machines and coordinate distributed commits.
One database, many machines
Replication keeps copies of data on several nodes for availability and read scaling. In primary-replica (leader-follower) setups, writes go to the primary and stream to replicas; synchronous replication waits for replicas before committing (no data loss on failover, higher latency), while asynchronous replication does not (faster, but recent commits can be lost on failover and replicas can serve stale reads). Partitioning (sharding) splits data across nodes by range or hash of a shard key, so each node handles part of the data and writes; good shard keys spread load evenly and keep related data together. When one transaction touches data on several nodes, a distributed commit protocol is needed. Two-phase commit (2PC): a coordinator asks all participants to prepare (vote yes after making their changes durable but not final), and if all vote yes it tells them to commit, otherwise to abort. 2PC guarantees atomicity but can block if the coordinator fails at the wrong moment, which is why many systems use consensus protocols such as Paxos or Raft, or avoid distributed transactions with designs such as sagas.
Two-phase commit, step by step
A transfer between accounts stored on two different nodes.
coordinator C, participants N1 (account A), N2 (account B)
phase 1 - prepare
C -> N1, N2 : PREPARE
N1 : debit A in a durable, uncommitted state -> votes YES
N2 : credit B in a durable, uncommitted state -> votes YES
phase 2 - commit
C : writes COMMIT to its log
C -> N1, N2 : COMMIT -> both make the change final and release locks
if any participant votes NO (or times out) -> C sends ABORT to all
blocking case: participants voted YES and C crashes before phase 2 ->
they must wait, holding locks, until C recoversA wedding ceremony
The officiant (coordinator) asks each partner "do you?" (prepare). Only if both say yes does the officiant declare them married (commit). If the officiant faints after both said yes, everyone waits in uncertainty until they recover.
Quick check: What is a known weakness of two-phase commit?
- It cannot guarantee atomicity
- Participants can block, holding locks, if the coordinator fails after they voted yes
- It only works on one machine
- It requires NoSQL databases
Answer
Participants can block, holding locks, if the coordinator fails after they voted yes — 2PC is atomic but blocking when the coordinator fails at a critical point.