skip to content

Core Distributed Theory

The theory that everything else rests on: replication, consensus, CAP, partitioning, clocks and ordering, failure handling, and distributed transactions. Expect questions here to be about why something is impossible as often as about how to build it.

part ofDistributed & scalable systemsoverview, primer and where to startread it →
on this pageshow

questions

64 · 11 sections

In a single-leader (primary/replica) replication setup, how do writes and reads flow through the system, and why would a team adopt this design over just running one database server?

level: juniorimportance: must knowfreq 75%
basics
~20 s

One server, the leader, is the only one allowed to accept writes. It records each change and sends a copy to the other servers (followers). Reads can be spread across the leader and followers. Teams do this to survive a server crash and to handle more read traffic than one machine could.

open as a page

What causes replication lag between a leader and its followers, and how would you design a system to guarantee 'read-your-writes' consistency for a user despite that lag?

level: middleimportance: must knowfreq 80%
basics
~20 s

Followers apply changes slightly after the leader because copying and applying takes time, especially under load. Read-your-writes means a user should always see their own recent changes - you guarantee it by making sure that right after someone writes, their next read goes to the leader or to a follower confirmed to be caught up, instead of a random possibly-stale replica.

open as a page

In a single-leader replication setup, what's the practical difference between synchronous and asynchronous follower replication, and what does each cost you?

level: middleimportance: must knowfreq 70%
basics
~20 s

Synchronous means the leader waits for a follower to confirm it got the write before telling the client 'done' - safer but slower. Asynchronous means the leader tells the client 'done' immediately and sends the copy in the background - faster but a crash can lose the last few writes.

open as a page

When would a team choose multi-leader replication over single-leader, and what fundamentally new problem does allowing more than one writable node introduce?

level: seniorimportance: must knowfreq 55%
basics
~20 s

Multi-leader means more than one server can accept writes, usually one per data center or region, so users write to a nearby server instead of one far away, and the system stays writable even if one region goes offline. The catch: two leaders can accept conflicting writes to the same piece of data at nearly the same time, and someone has to decide which one wins.

open as a page

In a leaderless (Dynamo-style) replication system where any of N replicas can accept a write, how do write quorums (W) and read quorums (R) work together to make reads see recent writes, and what do read repair and hinted handoff do?

level: seniorimportance: should knowfreq 45%
basics
~30 s

There's no single leader - a client writes to several replicas at once and only needs a certain number (W) to confirm before it's considered done; reads similarly query several replicas (R) and use the newest answer. If W+R is more than the total number of replicas, at least one replica in any read overlaps with one in the write, so a recent write is very likely seen. Read repair fixes replicas caught with stale data during a read; hinted handoff lets a temporarily unreachable replica's write be held by another node and delivered later.

open as a page

In a Raft or Paxos cluster of N nodes, a write is considered committed once it has been acknowledged by a quorum. Why does that quorum have to be a strict majority (more than N/2) rather than some fixed count like 2 nodes, regardless of cluster size?

level: juniorimportance: must knowfreq 70%
basics
~10 s

A majority is needed so any two quorums always share at least one node - that overlap is what stops two different decisions from being made at once.

open as a page

Walk through how a Raft leader replicates a new client write to its followers and decides when that entry is safely 'committed.'

level: middleimportance: must knowfreq 75%
basics
~10 s

The leader adds the write to its own list first, sends copies to the other servers, and once more than half of them have saved it, the leader tells everyone it's official.

open as a page

In the context of consensus protocols like Raft, what's the difference between a 'safety' property and a 'liveness' property, and give an example of each being satisfied even when the other is temporarily violated?

level: middleimportance: must knowfreq 55%
basics
~20 s

Safety means 'nothing bad ever happens' (like two leaders both thinking they won), and liveness means 'something good eventually happens' (like a write eventually succeeding). A system can stay safe while briefly failing to make progress.

open as a page

Suppose a network partition isolates the current Raft leader from the rest of the cluster, but that leader doesn't crash - it just can't reach anyone. The majority side elects a new leader. When the partition heals, how does the protocol prevent both leaders from being treated as valid, and what stops the old leader from having accepted writes in the meantime?

level: seniorimportance: must knowfreq 60%
basics
~20 s

Every election bumps a counter called a term. The old leader's term is now lower than the new one's, so once they reconnect, everyone - including the old leader - recognizes the higher term and steps down; and the old leader couldn't commit any writes alone anyway since it needs majority acks it never got.

open as a page

A distributed system needs multiple services to agree on whether a single business transaction succeeds or fails, atomically. When would you reach for a two-phase-commit-style atomic commit protocol instead of running the operation through a full consensus protocol like Raft or Paxos, and what do you give up either way?

level: seniorimportance: should knowfreq 45%
basics
~20 s

2PC is for getting several different services to agree on one all-or-nothing transaction right now; consensus (Raft/Paxos) is for keeping several copies of the same data in agreement over time, and it survives a coordinator crashing, while plain 2PC can freeze if its coordinator dies mid-transaction.

open as a page

In plain terms, what does the CAP theorem say a distributed database must sacrifice when a network partition occurs, and why can't a system just have all three of consistency, availability, and partition tolerance?

level: juniorimportance: must knowfreq 80%
basics
~20 s

When two halves of a system can't talk to each other, you must pick: either every request works but might return old/wrong data (available), or requests wait/fail until the data is provably correct (consistent). You can't have perfect answers and instant answers at once during that outage.

open as a page

PACELC extends the CAP theorem with a clause for when the network is NOT partitioned. What does the 'ELC' part of PACELC add, and how does it change a system's design compared to only thinking in CAP terms?

level: middleimportance: must knowfreq 55%
basics
~20 s

PACELC says: if there's a network split (P), pick availability or consistency (A or C) — that's CAP. But Else (E), even when the network is fine, you still must pick between faster answers (Latency) or more correct answers (Consistency), because waiting to check with other machines always costs time.

open as a page

Some engineers describe a system as 'CA' — consistent and available, with no partition tolerance needed. Why do most distributed-systems practitioners consider 'CA' not a real option for a system that operates over an actual network, and what is the person usually actually describing when they use that label?

level: middleimportance: must knowfreq 45%
basics
~20 s

Real networks break sometimes — cables get cut, switches fail, regions lose connection. A system that spans more than one machine over such a network WILL face a partition eventually, so it can't promise to just never deal with one. 'CA' usually really means 'we haven't seen a partition happen yet' or 'we treat the whole system as down if one does.'

open as a page

You're picking a data store for a distributed lock service used for leader election versus a data store for shopping-cart sessions on a high-traffic e-commerce site. Walk through why you'd lean CP for one and AP for the other, and name the concrete cost each choice imposes on that specific system during a network partition.

level: seniorimportance: must knowfreq 70%
basics
~20 s

For the lock service, two 'leaders' at once could break everything, so it's worth going offline rather than risk that (CP). For shopping carts, showing an old cart is annoying but rarely catastrophic, and going offline during checkout loses sales, so keep answering even with slightly stale data (AP).

open as a page

A wide-column store lets you configure the consistency level per read or write request — for example, requiring acknowledgment from just one replica versus requiring acknowledgment from a majority of replicas. How does this per-request knob let a team make different CAP and PACELC trade-offs for different operations within the SAME cluster, and what's the risk of misusing it?

level: seniorimportance: should knowfreq 50%
basics
~20 s

Instead of the whole database picking 'always fast' or 'always safe' once, each individual read or write can choose: talk to just one copy (fast, might be stale) or talk to most copies (slower, more sure to be correct). Different parts of the app can pick differently. The risk is picking 'fast' for something that actually needed to be 'safe' and getting wrong data where it matters.

open as a page

When you split a dataset across multiple nodes, what's the difference between range partitioning and hash partitioning, and what kind of query does each make fast or slow?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Range partitioning keeps keys in sorted order across nodes (A-M on node 1, N-Z on node 2), so range scans are fast but one busy range can overload a node. Hash partitioning scrambles keys by a hash function, spreading load evenly but killing range scans.

open as a page

In a hash-partitioned cluster, plain hash(key) mod N remaps almost every key when N (the node count) changes. Explain how consistent hashing arranges nodes and keys on a ring to avoid this, and why adding or removing one node only affects a small fraction of keys.

level: middleimportance: must knowfreq 75%
basics
~20 s

Consistent hashing places both nodes and keys on a circular number line (a ring) using the same hash function; each key belongs to the next node clockwise from it. Adding or removing a node only shifts the keys between it and its neighbor, not the whole dataset, unlike mod-N hashing where nearly everything gets reshuffled.

open as a page

In a partitioned datastore where the primary key determines which partition a row lives on, what's the difference between a local secondary index and a global secondary index on some other (non-partition-key) attribute, and what does each cost you on reads versus writes?

level: seniorimportance: must knowfreq 55%
basics
~20 s

A local secondary index lives on the same partition as the data it indexes, so writes stay on one partition but a query has to check every partition to be complete. A global secondary index is its own separate partitioned structure, so a query can hit just the index's partitions directly, but every write now has to update two different partitions.

open as a page

On a consistent-hashing ring with only a handful of physical nodes, why can data and load still end up badly uneven, and how do virtual nodes (multiple ring positions per physical node) fix that?

level: middleimportance: should knowfreq 60%
basics
~20 s

With few nodes, random ring placement can give one node a much bigger arc than another, so it gets more data and traffic. Virtual nodes give each physical machine many small positions scattered around the ring instead of one big one, averaging out into roughly equal shares.

open as a page

When nodes are added to or removed from a partitioned cluster, walk through how rebalancing actually moves data. What are the main strategies (e.g., a fixed number of partitions per node vs. dynamically splitting partitions by size), and what can go wrong during the process?

level: seniorimportance: should knowfreq 65%
basics
~20 s

Rebalancing is moving data between nodes so each one holds a fair share after the cluster changes size. One approach creates many more partitions than nodes up front and just reassigns whole partitions to nodes as they join/leave; another splits/merges partitions dynamically as they grow or shrink. Done badly, it can overload the network or serve stale/missing data mid-move.

open as a page

Two servers in different data centers each timestamp an event using their local system clock (wall-clock/time-of-day). Why can't you safely assume that whichever timestamp is numerically smaller happened first in real time?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Computer clocks aren't perfectly synced - they drift apart between corrections and can even jump backward during a correction, so timestamps from two different machines can't be trusted to show which event really happened first.

open as a page

A service measures how long an operation took by calling a wall-clock/time-of-day API (e.g. the equivalent of `System.currentTimeMillis()`) before and after the operation, then subtracting. Under what circumstances can this produce a negative or wildly wrong duration, and what's the correct fix?

level: middleimportance: must knowfreq 55%
basics
~20 s

If the system's clock gets adjusted backward (say by a time-sync correction) while you're timing something, the 'end' reading can look earlier than the 'start' reading, giving a negative or nonsense duration. Use a monotonic clock instead - one guaranteed to only ever move forward.

open as a page

When a server synchronizes its clock via NTP (Network Time Protocol), walk through how NTP estimates the clock offset and round-trip delay from a time server, and explain why the resulting accuracy is fundamentally limited by network conditions rather than by protocol design.

level: middleimportance: must knowfreq 65%
basics
~20 s

NTP asks a time server what time it is, times how long the round trip took, and assumes the trip there and back took equally long to estimate the network delay and correct the local clock. If the trip isn't actually symmetric, the correction is a little off.

open as a page

Google's TrueTime API, used inside the Spanner database, doesn't return a single timestamp for 'now' - it returns an interval [earliest, latest]. Explain the mechanism behind this design, including the role GPS receivers and atomic clocks play, and why returning an interval is more useful than a system trying to report one perfectly accurate timestamp.

level: seniorimportance: must knowfreq 40%
basics
~20 s

Instead of pretending to know the exact time, TrueTime admits 'the real time is somewhere in this small window,' using GPS satellites and atomic clocks spread across data centers to keep that window tiny. This gives the database a guaranteed, honest bound to reason with instead of a false, precise-looking number.

open as a page

In Google Spanner, a read-write transaction's commit protocol includes a 'commit-wait' step where the coordinator delays making the transaction's writes visible until a certain point. Describe what commit-wait actually waits for and why it's necessary to guarantee external consistency - meaning transactions appear to execute in an order consistent with real, wall-clock time.

level: seniorimportance: should knowfreq 25%
basics
~20 s

Spanner picks a commit timestamp for a transaction, then literally pauses before letting anyone see the result, until it's sure real-world clocks everywhere have caught up past that timestamp. This guarantees that anything starting after the commit will see it, at the cost of adding a little delay to every write.

open as a page

In distributed systems, engineers talk about crash-stop, crash-recovery, omission, and Byzantine failure models. What does each model assume about how a faulty node can misbehave, and why does the choice of model matter when designing a fault-tolerant system?

level: juniorimportance: must knowfreq 65%
basics
~20 s

A failure model is an assumption about how something can break. Crash-stop: a node dies and never comes back. Crash-recovery: it dies but can restart. Omission: messages get lost. Byzantine: a node can lie or act maliciously. Stronger assumptions need more expensive protection.

open as a page

A replicated key-value store has N=5 replicas. A team configures writes to require acknowledgment from W=3 replicas and reads to query R=3 replicas. Why does this W+R>N configuration let the system mask node failures while still returning consistent reads, and what happens if the team instead sets W=2, R=2?

level: middleimportance: must knowfreq 62%
basics
~20 s

With W+R>N, any read set and any write set must share at least one replica, so a read always sees the latest write even if some nodes are down or slow. If W+R is not greater than N, like W=2,R=2 on N=5, reads and writes might miss each other entirely and return stale data.

open as a page

A team is deciding between synchronous and asynchronous replication for a database that must tolerate a primary node failure. What availability and durability guarantees does each give, and what specific failure mode does asynchronous replication risk that synchronous replication avoids?

level: middleimportance: must knowfreq 68%
basics
~20 s

Synchronous replication waits for a copy to confirm the write before telling the client it succeeded, so no acknowledged write is ever lost, but it's slower and can stall if a replica is unreachable. Asynchronous replication confirms instantly and copies data in the background, so it's fast, but a crash right after can lose the last few writes.

open as a page

A distributed lock service grants a lease to a client believed to hold exclusive access to a shared resource, such as a storage volume. The client experiences a long GC pause, its lease expires, and a second client acquires the lease and starts writing. The first client then resumes and, unaware its lease expired, also writes. How does fencing prevent this from corrupting the shared resource, and why isn't simply checking the lease is still valid on the client side sufficient?

level: seniorimportance: should knowfreq 50%
basics
~20 s

Fencing gives each lease a rising number; the resource itself refuses any write whose number is older than the newest it has seen. That way even a confused old client that thinks it still owns the lock physically can't write, because the resource, not the confused client, enforces the check.

open as a page

The FLP impossibility result states that no deterministic consensus protocol can guarantee both safety and termination in a fully asynchronous system where even a single process may crash. Given that systems like Raft and Paxos are used in production to reach consensus every day, how do real systems reconcile their existence with this theoretical impossibility?

level: principalimportance: should knowfreq 35%
basics
~20 s

FLP proves that, in theory, a perfectly asynchronous network can always delay messages just enough to stall a consensus algorithm forever. Real systems dodge this by accepting that consensus might occasionally stall for a while, never violating correctness, rather than promising it will always finish quickly - and by using timeouts that work well enough in practice even though they aren't a formal guarantee.

open as a page

In a distributed system with no shared physical clock, why can't you just compare timestamps from different machines to figure out which of two events happened first?

level: juniorimportance: must knowfreq 65%
basics
~20 s

Different computers' internal clocks are never perfectly synchronized, so their timestamps can disagree about event order even when one event actually caused the other. You need a way to order events based on what actually influenced what, not on unreliable clock readings.

open as a page

Describe the update rules for a Lamport timestamp counter, and explain why two events having comparable Lamport timestamps (one number less than the other) does NOT guarantee that the smaller one happened-before the larger one.

level: middleimportance: must knowfreq 70%
basics
~20 s

Each process keeps a counter it bumps for every local event, and includes it in messages; a receiver sets its counter to max(own, received)+1. Bigger numbers don't prove causation -- two unrelated events on different machines can just happen to get different counter values.

open as a page

A system uses vector clocks with one integer slot per process. Walk through the update rules for incrementing and merging a vector clock, and explain precisely how comparing two vector clocks lets you detect that two events are concurrent rather than causally ordered.

level: middleimportance: must knowfreq 65%
basics
~20 s

Each process keeps a list of counters, one per process, not just one number. On its own events it bumps its own slot; on receiving a message it takes the slot-by-slot maximum with the sender's vector, then bumps its own slot. If neither vector is fully >= the other, the events are concurrent -- that's the extra power a single Lamport number can't give you.

open as a page

Explain what a Hybrid Logical Clock (HLC) adds on top of a plain Lamport timestamp, how its (physical-time, logical-counter) pair is updated on local events and message receipt, and what problem this hybrid design solves that neither pure physical clocks nor pure Lamport clocks solve alone.

level: seniorimportance: should knowfreq 45%
basics
~20 s

An HLC pairs a physical wall-clock reading with a small logical counter that only increments when needed to break ties or preserve causality. That gives timestamps that stay close to real time (useful for humans and TTLs) while still guaranteeing that causally related events get increasing values, which plain physical clocks alone can't guarantee under clock skew.

open as a page

In a leaderless replicated key-value store like Amazon's Dynamo, replicas track causality per-replica rather than per-client-request. Explain the difference between a version vector (keyed by replica) and a full vector clock (keyed by every process/event-generating actor), and why the store picks the former.

level: seniorimportance: should knowfreq 55%
basics
~20 s

A version vector has one counter per data replica/server, updated when a replica handles a write, rather than one counter per every client or process that ever touches the system. It's a cheaper, coarser-grained cousin of a full vector clock -- good enough to detect conflicting replica states without needing a slot for every client that ever connects.

open as a page

In a cluster of identical worker nodes that each have a unique numeric ID, why do distributed systems need a 'leader election' process at all, and how does the classic Bully algorithm decide who becomes leader?

level: juniorimportance: must knowfreq 55%
basics
~20 s

Letting two computers both act as 'the boss' at once causes conflicts. Leader election is how a group agrees on exactly one boss. The Bully algorithm is simple: whichever live computer has the highest ID number always wins.

open as a page

In Raft, nodes use randomized election timeouts, terms, and majority votes to pick a leader. Walk through exactly how a follower becomes a candidate and wins an election, and explain why the timeout is randomized rather than fixed.

level: middleimportance: must knowfreq 70%
basics
~20 s

Every node waits a random amount of time to hear from the leader. Whoever times out first asks everyone to vote for it. Getting votes from more than half the nodes wins. The random wait stops everyone timing out together.

open as a page

What does it mean for a node to hold leadership via a time-bound 'lease' rather than an indefinite election result, how does lease renewal work, and what production failure mode does clock skew or a stop-the-world pause introduce?

level: seniorimportance: must knowfreq 45%
basics
~20 s

A lease is leadership with an expiration timer, like a rented spot you keep paying for. Renew before it runs out or lose leadership. Danger: a frozen leader might think its lease is still valid after someone else took over.

open as a page

Beyond electing a new leader correctly, what concrete mechanisms stop a still-running, previously-elected leader from continuing to write to shared storage after a network partition has caused a new leader to be chosen, and why is preventing the election alone not enough?

level: seniorimportance: must knowfreq 40%
basics
~20 s

Even after everyone agrees on a new leader, the old leader may not know it's been replaced and could keep writing. Fix: storage refuses writes from anyone but the current leader, via a growing 'ticket number' that rejects old ones.

open as a page

How does the Ring leader-election algorithm, where each node only ever talks to its logical successor around a virtual ring, work mechanically, and what does it trade off against the Bully algorithm's approach of contacting every higher-ID node directly?

level: middleimportance: should knowfreq 25%
basics
~20 s

Nodes form a logical circle, each only talking to the next node. A node's ID gets passed around, each node swapping in its own ID if it's bigger, until the message returns carrying the biggest ID - that node wins.

open as a page

A distributed transaction coordinator needs to make sure a payment write to Database A and an inventory write to Database B either both happen or neither happens. Walk through how the two-phase commit (2PC) protocol achieves this, phase by phase.

level: juniorimportance: must knowfreq 75%
basics
~10 s

2PC asks everyone 'can you commit?' first (prepare phase). Only if all say yes does it tell them to actually commit (commit phase). If anyone says no, everyone rolls back.

open as a page

In a two-phase commit transaction, every participant has replied VOTE-COMMIT and is holding its locks, waiting for the coordinator's final instruction — but the coordinator crashes before sending COMMIT or ABORT to anyone. Why can't the participants just make a decision themselves, and what state does this leave the system in?

level: middleimportance: must knowfreq 65%
basics
~20 s

The participants can't tell if the coordinator had already decided to commit or abort before it died, so guessing wrong could break consistency. They're stuck holding locks — 'blocked' — until the coordinator comes back or someone intervenes.

open as a page

A team is designing a checkout flow that needs to debit a payment account, decrement inventory, and create a shipping order across three separate services. They're deciding between wrapping the whole thing in a two-phase-commit (2PC) distributed transaction versus using a Saga (a sequence of local transactions with compensating actions to undo earlier steps if a later one fails). What factors should drive that decision, and what does each approach give up?

level: seniorimportance: must knowfreq 60%
basics
~20 s

2PC keeps everything locked until all three services agree, so it's always consistent but slow and fragile if one service is down. A Saga lets each step commit right away and fixes mistakes afterward with compensating actions, so it's faster and more resilient but briefly shows an inconsistent state.

open as a page

An operations team notices that after a brief network blip, several database connections show 'in-doubt' or 'prepared' transactions that never resolved, and those transactions are holding row locks that are blocking other queries. In the context of two-phase commit / XA transactions, what causes this, and what are the operator's options to resolve it?

level: middleimportance: should knowfreq 45%
basics
~20 s

A brief network blip cut the connection between the coordinator and a database right after it said 'ready to commit', so the database is stuck waiting and holding locks. An operator can wait for the coordinator to reconnect and resolve it automatically, or manually force a commit/rollback if it's urgent.

open as a page

Three-phase commit (3PC) adds a 'pre-commit' phase between the initial prepare vote and the final commit, specifically to let surviving participants recover without waiting indefinitely for a crashed coordinator. Explain how that extra phase helps under a simple crash-stop failure model, and why 3PC still doesn't solve the blocking problem once you allow network partitions.

level: seniorimportance: should knowfreq 35%
basics
~20 s

3PC adds a middle step so participants know 'everyone agreed to commit' before actually committing, letting them safely finish on their own if the coordinator dies. But if the network splits instead of the coordinator crashing, participants on different sides can still disagree, so it still gets stuck.

open as a page

In distributed systems, what is the difference between at-most-once, at-least-once, and exactly-once message delivery, and which of these can actually be guaranteed over an unreliable network?

level: juniorimportance: must knowfreq 85%
basics
~20 s

At-most-once means a message might get lost but never duplicated. At-least-once means it might arrive more than once but never lost. Exactly-once (arrives precisely one time) can't be truly guaranteed when networks can drop or delay messages - only faked using retries plus deduplication.

open as a page

A client sends POST /charges to bill a customer's card. The server processes the charge successfully, but the network drops before the response reaches the client, so the client's HTTP library times out and automatically retries the identical request. How does an idempotency key prevent the customer from being charged twice?

level: middleimportance: must knowfreq 90%
basics
~20 s

The client attaches a unique key to the request. The server remembers which keys it already handled and what it returned, so when the same key shows up again, it just replays the original result instead of charging the card a second time.

open as a page

Two identical retry requests carrying the same client-generated idempotency key arrive at two different, stateless replicas of an API server within a few milliseconds of each other - before either has finished writing to the shared idempotency store. What race condition can occur, and what pattern prevents it?

level: seniorimportance: must knowfreq 55%
basics
~20 s

Both replicas might check the store, see nothing there yet, and both go ahead and process the request - so the action (like a charge) happens twice. The fix is to make 'claiming' the key an atomic, all-or-nothing step that only one replica can win.

open as a page

A background worker consumes events from an at-least-once message queue and updates a search index for each event. Some events get redelivered after worker restarts. How would you use a request-deduplication store to make this consumer effectively idempotent, and what has to go into that store?

level: middleimportance: should knowfreq 65%
basics
~20 s

The worker keeps a record of every event ID it has already handled. Before processing a new event, it checks that record; if the ID is already there, it skips the event instead of applying it again.

open as a page

An operation like 'increment a user's loyalty-points balance by 10' or 'send a welcome email' is not naturally idempotent - running it twice produces a different (or duplicated) outcome than running it once. Given that the underlying operation can be delivered more than once, what general design techniques make it safe to apply repeatedly?

level: seniorimportance: should knowfreq 50%
basics
~20 s

Instead of blindly repeating the action, redesign it so repeats do nothing extra: turn 'add 10 points' into 'record this specific +10 transaction with an ID', or check a 'have I already emailed this user for this event' flag before sending.

open as a page

A cluster node pings its peers every second and marks a peer 'dead' if it misses 5 heartbeats in a row (5 seconds of silence). What is the basic trade-off this fixed-timeout heartbeat approach faces when choosing the timeout value, and why can't a single value be 'correct' for all conditions?

level: juniorimportance: must knowfreq 65%
basics
~10 s

A short timeout catches crashes fast but often wrongly declares a slow-but-alive node dead. A long timeout avoids false alarms but takes longer to notice a real crash.

open as a page

The phi-accrual failure detector (used in Cassandra and Akka Cluster) doesn't output a binary alive/dead like a fixed heartbeat timeout — it outputs a continuous suspicion level called phi. What does phi represent, and why is that more useful than a single global timeout?

level: middleimportance: must knowfreq 60%
basics
~10 s

Phi is a number that grows the longer a node stays silent past its normal heartbeat rhythm; the app picks how high phi must get before acting, instead of one hard-coded timeout.

open as a page

In the SWIM (Scalable Weakly-consistent Infection-style Membership) protocol, when node A directly pings node B and gets no response, A doesn't immediately declare B dead — it asks a few other nodes to indirectly probe B on its behalf. Why does SWIM add this indirect-probe step before declaring a peer suspect, and how does it help SWIM scale to large clusters?

level: seniorimportance: must knowfreq 55%
basics
~20 s

Because A's own path to B might just be a bad connection, not B actually crashing — asking others to check too avoids blaming B for a problem that's really just between A and B.

open as a page

A gossip/epidemic protocol spreads membership updates by having each node periodically pick a few random peers and exchange state. Comparing 'push' gossip (a node sends its state to random peers), 'pull' gossip (a node asks random peers for their state), and 'push-pull' (both directions in one exchange), what's the practical trade-off between them in terms of convergence speed and network cost, and where does anti-entropy fit in?

level: seniorimportance: should knowfreq 40%
basics
~20 s

Push spreads new information fast at first but wastes messages once most nodes already know it; pull is good at catching stragglers but slow to start; push-pull combines both and converges fastest but costs the most bandwidth per round. Anti-entropy is a slower background pass that repairs anything the fast path missed.

open as a page

Because gossip-based membership propagation is inherently asynchronous, different nodes in the same cluster can briefly hold different views of 'who is alive' — node A might already consider node C dead while node B still considers C alive. What are the practical implications of this eventually-consistent membership view, and how do production systems mitigate the risk of inconsistent liveness decisions?

level: principalimportance: should knowfreq 28%
basics
~20 s

Since gossip takes time to spread, nodes can briefly disagree about who's alive or dead; systems handle this by not trusting any single node's opinion alone and requiring agreement, refutation, or an explicit confirmation step before acting.

open as a page