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, like pg_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.

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.
    • 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
  • 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