F1 Leader Election & Distributed Coordination
Loading learning experience...
Lecture transcript
Read the narration for F1: Leader Election & Distributed Coordination
From TAO's graph to: who coordinates changes?
Dr. Wei: Last time, we looked inside TAO: objects and associations are the two primitives, and the hard part is serving a huge graph reliably at scale.
Dr. Wei: Today we add the missing ingredient: when many machines store and cache that graph, who gets to decide shard assignment, config updates, and rebalancing without corrupting state?
Sam: So this is the part where one node has to be 'in charge', but we can't trust any node to stay alive?
Dr. Wei: This bullet is why leader election shows up everywhere: some actions must be decided exactly once, otherwise you create conflicting realities.
Dr. Wei: Our goal is not 'pick a leader once'. It is: keep a single leader property even while nodes crash and networks split.
Sam: And if we fail, we get two leaders and chaos?
Dr. Wei: In interviews, you can go either way: embed consensus inside the system, or rely on an external CP coordinator that everyone trusts.
Dr. Wei: Next, we will make split-brain concrete, then we will build the two standard fixes: majority quorum and fencing.
Why leaders: exactly-once decisions in a messy world
Dr. Wei: Leader election sounds academic until you ask one simple question: who assigns shards when machines are failing mid-flight?
Dr. Wei: By the end, you will be able to justify a safe coordination choice and explain how it avoids split-brain.
Dr. Wei: We will start from simple failure stories, then connect them to quorum voting, ephemeral nodes, and fencing tokens.
Sam: I get the need for a leader, but I do not see why 'two leaders' is so catastrophic.
Dr. Wei: These are the operations where you want a single source of truth: they change the system's structure, not just data values.
Dr. Wei: Imagine two schedulers both think a job is unassigned: you get duplicate side effects. For rebalances, you can move the same shard twice in different directions.
Dr. Wei: Notice the priority order: safety first. Many real coordinators will deliberately stop making progress rather than risk two leaders.
Sam: So we are okay with a pause, but not okay with conflicting commands?
Dr. Wei: Here is what the interviewer is looking for: you state the failure model explicitly, then you pick a mechanism that is correct under that model.
Split-brain: how $2$ leaders happen
Dr. Wei: Let us build the simplest split-brain story with three nodes, because it forces us to reason about partitions, not just crashes.
Dr. Wei: Start with a normal state: A is issuing coordination decisions and B and C are following.
Dr. Wei: The key detail is that A is alive but cut off. From A's perspective, nothing proved it is not leader anymore.
Dr. Wei: Meanwhile B and C see no heartbeat from A, so they do what a naive system would: elect a new leader inside their partition.
Sam: So now A and B both think they are leader, and neither is actually lying.
Dr. Wei: This is where corruption happens: A may assign shard X to one host, B assigns shard X to a different host, and now clients see nondeterministic routing and data loss.
Dr. Wei: The classic fix is quorum: only a partition with a majority can elect a leader. With three nodes, only B plus C can form that majority; A alone cannot.
Quorums: why majority gives uniqueness
Dr. Wei: Now we will turn that intuition into a one-line guarantee: two majorities must overlap, so two leaders cannot both be valid.
Sam: Overlap is doing all the work here, right?
Dr. Wei: This is the quorum size: more than half the nodes. It is the smallest size that forces any two quorums to share at least one node.
Dr. Wei: If two different leaders both claim legitimacy, each must have gathered a majority. But two majorities share a voter, and that voter cannot validly vote for both at the same time.
Dr. Wei: This is the CAP implication: if the network splits into two halves, neither side has a majority, so the system prefers to stop coordinating rather than risk split-brain.
Sam: So availability is sacrificed specifically for coordination operations.
Dr. Wei: Strong answers separate concerns: data reads may still work from caches, but metadata changes like shard moves should freeze until a leader is safely elected.
Raft election: randomized timeouts $+$ majority vote
Dr. Wei: Raft is a common way to implement quorum leadership: it makes elections understandable and avoids constant leader churn.
Dr. Wei: Most of the time nodes are followers. A candidate state is just a temporary mode while trying to become leader.
Dr. Wei: Random timeouts reduce collisions: if every follower timed out at the same instant, everyone would start an election and no one would win.
Sam: So randomness is not a performance trick, it is a correctness stability trick.
Dr. Wei: When a node starts an election, it bumps a term number and asks peers to vote. The term is like an epoch: it lets the cluster reject stale leadership claims.
Dr. Wei: Once the candidate gets a majority, it becomes leader and immediately sends periodic heartbeats. The point is to keep followers calm so they do not start competing elections.
Dr. Wei: This is the key safety guardrail: if an old leader comes back from a partition, the term comparison forces it to step down when it sees a newer term.
Raft replication: leader is not just a title
Dr. Wei: In consensus, leadership matters because the leader orders operations, and everyone applies the same ordered log.
Sam: So election alone is not enough; we also need agreement on the sequence of changes.
Dr. Wei: The log is the leader's proposal of history. Every coordination decision becomes a log entry with an index.
Dr. Wei: Followers do not invent new entries; they only accept entries from the leader for the current term, which keeps the history linear.
Dr. Wei: This is the same quorum idea again: majority acknowledgment means that future leaders must overlap with at least one node that knows this entry exists.
Dr. Wei: This is how coordination becomes useful: the log is not the goal. The goal is that everyone reaches the same state by applying the same sequence of operations.
Dr. Wei: This is the 'no time travel' property: once the cluster agrees an operation happened, leader changes cannot erase it, even if the old leader dies.
ZooKeeper-style coordination: tiny CP database for agreement
Dr. Wei: The other common pattern is: do not embed consensus everywhere; instead, use a small, highly reliable coordination service that provides primitives.
Sam: So ZooKeeper is basically a tiny, ultra-reliable database that everyone uses to agree on who is in charge?
Dr. Wei: ZooKeeper is optimized for coordination metadata: membership lists, locks, leader pointers. If you try to store high throughput business data, you will overload it.
Dr. Wei: These three types are the building blocks. Persistent is like a normal row. Ephemeral is tied to a live session. Sequential gives a total order by auto-incrementing names.
Dr. Wei: Watches are how you avoid polling. You set a watch, and when membership changes, you wake up and recompute assignments.
Dr. Wei: This is why it is safe for leader election: operations appear in a single global order, even during failures. The cost is that a partition can make it unavailable.
ZooKeeper leader election recipe: ephemeral sequential nodes
Dr. Wei: Now we will do the classic recipe you can describe in interviews step by step, including herd avoidance and fencing.
Dr. Wei: Each candidate creates a node like node dash 0007. Ephemeral means it disappears if the session dies; sequential means the service assigns an increasing suffix.
Dr. Wei: This gives you a deterministic leader with no extra voting protocol: the order is already provided by the coordination service.
Dr. Wei: A common pitfall is the thundering herd: if everyone watches the leader, then one leader death wakes up ten thousand clients. Watching only your predecessor localizes notifications.
Sam: So only one follower wakes up on each leader change, then it becomes leader.
Dr. Wei: This is the elegance of ephemeral: session liveness is the lease. When the lease ends, the coordination record is cleaned up automatically.
Dr. Wei: And this is the IC6 detail: election alone is not enough. You attach a monotonically increasing token to every privileged write so storage can reject stale leaders.
Distributed locks: correctness gotcha and fencing fix
Dr. Wei: Locks are the simplest place to see the hardest bug: a client thinks it holds the lock, but the system disagrees after a session hiccup.
Sam: This is the 'I kept running, but my lease expired' problem?
Dr. Wei: This is identical to leader election, just scoped per resource r. You turn ordering into mutual exclusion.
Dr. Wei: Again, predecessor watches avoid waking everyone up. Only the next-in-line client gets the signal to retry.
Dr. Wei: If the client is paused by GC or a network blip, it might miss heartbeats. ZooKeeper deletes its ephemeral lock, a new client acquires it, but the old client continues executing critical-section code.
Dr. Wei: The lock service cannot stop a buggy or slow client from sending writes. So the storage layer must enforce: only the highest token is allowed to mutate state.
Sam: So the token is like a version number for authority, not just for debugging.
Dr. Wei: If you mention distributed locks in a Meta interview, you should mention fencing within the next sentence. That is a strong senior correctness signal.
Applied example: Kafka controller election and broker membership
Dr. Wei: Let us ground this in a real system design story: Kafka needs a single controller to assign partitions and handle rebalancing.
Dr. Wei: The controller is the metadata leader. If you had two controllers, you could assign the same partition to two leaders and clients would observe divergent writes.
Dr. Wei: Brokers create ephemeral records so the controller can treat disappearance as failure. This avoids building custom failure detection logic everywhere.
Dr. Wei: Watches turn metadata into a publish-subscribe channel: the controller updates assignment, and brokers receive notifications to start or stop serving partitions.
Dr. Wei: Many systems are moving to embedding consensus so you reduce dependencies and operational complexity. But you still pay the same quorum and CP tradeoffs.
Failure modes checklist: what breaks, what heals
Dr. Wei: Interviewers love a clean failure walkthrough: detection, election, fencing, and recovery of in-flight operations.
Sam: Can we do it as a checklist I can reuse in designs?
Dr. Wei: First you need a signal that leadership might be gone. In Raft it is missing heartbeats; in ZooKeeper it is the session lease ending.
Dr. Wei: Then you enforce quorum. Under a partition, you prefer to block rather than create two decision makers.
Dr. Wei: This handles the nasty case: an old leader is alive and still sending requests. Tokens let downstream systems reject stale authority even if the old leader is confused.
Dr. Wei: Finally, you need clean-up. With consensus you replay a committed log. With coordination metadata, you recompute assignments deterministically and make operations idempotent so retries are safe.
Sam: That helps: detect, elect, fence, recover.
Exit ticket: reason about a partition and show safety
Dr. Wei: Let us lock in the two key skills: quorum reasoning and fencing reasoning.
Dr. Wei: Compute quorum: with five nodes, you need three votes. In a two-and-three partition, only the three-node side can possibly collect three votes, so only that side can elect a leader.
Sam: So the two-node side must stay follower forever until the partition heals, even if it is serving reads.
Dr. Wei: Because election does not physically stop the old leader process. If it comes back and sends writes, the token is what lets storage prove it is stale and reject it, preventing split-brain side effects.
Dr. Wei: If you can say quorum, CP tradeoff, and fencing clearly, you will sound like a safe systems designer in interviews.
Thank you for watching!
Thanks for watching. Subscribe and share if you found this useful—see you next time!