ddia/ch06.md - DDIA
Single-leader replication
- Synchronous replication: Easy. Client can read from replicas
- Impractical. If all replicas are down, we can't operate.
- Sometimes, in a sense that 1 is sync, other async.
- Problem: data might be lost, if leader dies immediately after this replication. If the new leader doesn't have updated data, data is lost.
- Sometimes, in a sense that you sync to quorum, other async
- Solve the problem above.
- Asynchronous replication: Stale data is unavoidable. Possible solutions to make it bounded and
Sharded DB:
- If DB is sharded, each shard has 1 leader
How to add a new follower:
- Set up a follower node.
- Take snapshot of leader's DB.
- Roll out a snapshot of the leader's DB on a follower node.
- Follower connects to the leader and pulls the delta using the latest sequence number.
Object storage
Using object storage as a shared database. Benefits:
- Cheap storage
- Good replication
- Can use CAS If you have 3 different services and need to elect a leader, it's hard to do. Single coordination point with CAS support helps.
- Single storing place for many sources (usecase for Data Warehouses)
Drawbacks:
- Expensive requests => need to batch
- Slower
- No update feature. Objects are immutable.
Handling outages
- Follower failure: store log of processed operations, then just reboot and pull the delta from the leader. If the follower is far behind, set it up as a new follower (pull the snapshot).
- Leader failure (failover): either manual or automatic. Consider automatic:
- Determine that the leader has failed: timeout (health check).
- Choose a new leader.
- Either by election => consensus
- Or, other entity chooses it (based on highest WAL seq).
- Reconfigure the system to use the new leader.
- Redirect clients to a new leader.
- Use reverse proxy, or DNS, or VIP.
- Or, store leader address in a shared store (
etcd, ZooKeeper). Update the path there.
- Let other followers know the leader has changed
- In consensus-based algos, it's a part of the protocol
- Without consensus, update the shared store (e.g.
etcd, or other shared monitor, likepg_autoctl)
- Let the old leader know that it's now a follower
- After restart it contacts monitor to get its current role.
- What to do in network partition?
- Guard: if leader can't reach monitor => stop accepting writes.
- Redirect clients to a new leader.
Common problems of failovers
- New leader/old follower might not have the latest writes => Newer writes from new follower/old leader are discarded.
- This can become a problem if other systems need to be coordinated with the database contents.
Replication logs implementation
- Statement-based replication
- Idea: log of evens (insert, update, delete)
- Problems:
- Non-deterministic expressions (
rand(),now()). - Conditioned expressions (using
where). DBs do not necessarily have the same state. - Side effects => can diverge between replicas.
- Non-deterministic expressions (
- Solution: require transactions to be deterministic; expand expressions.
- WAL replication
- Idea: use DB's WAL
- Problems:
- Tight coupling to specific storage engine => if storage format is changed, we have a problem => leader and followers MUST have the same DB version => can't do zero downtime upgrades when storage format changes
- Solution: next method
- Row-based replication
- Idea: store logical changes - added rows (with values), deleted rows (their ids), updated rows (delta). This approach decouples representation from storage engine.
- Benefits:
- Backwards compatible => upgrade to new version with 0 downtime.
Mitigation of replication lag
- Read-your-writes consistency
- Idea: after user writes, they should see their write immediately
- How:
- Read from leader (for a period T after write = bounded staleness)
- Variation: session stickiness. read from node where write happened
- Versioning: remember the commit position on write and send it together with the read request.
- Quorum reads
- Client-side write cache
- Read from leader (for a period T after write = bounded staleness)
- Monotonic reads
- Idea: can't go back in time = we can't receive older data on repeated read.
- How:
- Versioning (see above).
- Read from only 1 replica (for example, based on user hash ID)
- Consistent prefix reads
- Idea: causal events must be retrieved in logical sequence.
- Important: weaker than causal consistency.
- Because causal consistency guarantees that if-A-then-B, we observe both in the right order, while consistent prefix only guarantees we don't see A after B (if-B-then-A).
- Problem: particularly present in sharded databases.
- Solution:
- Causal writes go to 1 shard / single leader per partition.
- Keep causal version between read requests
- Log-based systems