Consistency, Consensus and Distributed Transactions
When data lives on several machines, three questions dominate: which copy is the truth, how do machines agree on anything when messages can be lost or delayed, and how do you make a change that spans more than one machine all-or-nothing. This chapter gives you the vocabulary and the reasoning for each.
1. Why distributed systems are hard
Three facts cause nearly every difficulty:
- Networks fail. Messages are lost, delayed, duplicated and reordered.
- Clocks disagree. Machine clocks drift, so "which event happened first" cannot be read off timestamps.
- Partial failure. A request may time out, and you cannot tell whether the other side did the work or never received it.
The last one is the most important in practice. If a payment call times out, the payment may or may not have happened. Everything in this chapter exists to cope with that uncertainty.
2. Transactions and isolation
An ACID transaction gives atomicity (all or nothing), consistency (invariants hold), isolation (concurrent transactions do not interfere) and durability (committed data survives crashes).
Isolation is a spectrum, and the anomalies each level allows are a standard question:
| Level | Prevents | Still allows |
|---|---|---|
| Read uncommitted | Nothing much | Dirty reads |
| Read committed | Dirty reads | Non-repeatable reads, phantoms, write skew |
| Repeatable read / snapshot | Non-repeatable reads | Phantoms (varies), write skew |
| Serializable | All of these | Lower throughput |
- Dirty read: you read another transaction's uncommitted change.
- Non-repeatable read: the same row changes between two reads in your transaction.
- Phantom: a repeated query returns new rows.
- Write skew: two transactions read overlapping data, each updates a different row, and together they break an invariant (two doctors both go off call because each saw the other still on call). Snapshot isolation allows it, and only serializable prevents it.
Be ready to say that a lost update, where two read-modify-write cycles overwrite each other, is fixed by atomic operations, row locks, or compare-and-set with a version column.
3. Optimistic and pessimistic concurrency
Pessimistic control takes locks before touching data. It is safe under high contention and risks deadlock and waiting.
Optimistic control assumes conflicts are rare: read a version, do the work, and write only if the version is unchanged, otherwise retry.
UPDATE inventory SET qty = qty - 1, version = version + 1
WHERE sku = 'A1' AND version = 7 AND qty > 0;
-- if 0 rows updated, someone else got there first: retry or fail
Optimistic control wins when contention is low. Under heavy contention on one row, such as a flash-sale item, retries pile up, and you want a queue or a single writer for that item.
4. Idempotency: making retries safe
Since a timed-out request may have succeeded, clients retry, and retries must not do the work twice. An operation is idempotent if repeating it has the same effect as doing it once.
Natural idempotency: "set the status to shipped". Not idempotent: "add 10 to the balance".
For operations that are not naturally idempotent, use an idempotency key: the client generates a unique identifier per logical request. The server records the key with the result, and a repeat returns the stored result instead of executing again.
POST /payments Idempotency-Key: 7f3a... { amount: 500 }
The server stores the key and the response in the same transaction as the payment itself, so a crash cannot leave one without the other. This is the standard answer to "how do you avoid double charging".
Delivery guarantees for messages: at-most-once (may lose), at-least-once (may duplicate), and exactly-once. True exactly-once delivery over an unreliable network is not achievable in general. What systems provide is at-least-once delivery plus idempotent processing, which yields exactly-once effects. Say it that way.
5. Time, order and clocks
Wall clocks are unreliable for ordering events across machines. Two alternatives:
- Logical clocks (Lamport). Each machine keeps a counter, increments it on every event and sends it with messages, taking the maximum on receipt. If event A causes B, A's timestamp is smaller. The reverse does not hold.
- Vector clocks record a counter per machine and can tell whether two events are ordered or concurrent, which is needed to detect write conflicts.
Some systems bound clock uncertainty with specialised hardware and wait out the uncertainty before committing. You do not need the details, only the idea: global ordering can be bought with latency.
6. Quorums
In a replicated store with replicas, a write is acknowledged by replicas and a read consults . If
then every read set overlaps every write set, so a read sees the latest acknowledged write. With , choosing and tolerates one failed replica for both reads and writes. Choosing , is fast and eventually consistent. Choosing makes writes fail if any replica is down.
Quorums are tunable per request, which is how some stores let you trade latency for consistency.
<!--fig:quorum-->7. Consensus and leader election
Consensus is getting several machines to agree on a value despite failures. It underlies leader election, configuration stores and replicated logs. The well-known algorithms are Paxos and Raft; Raft was designed to be easier to understand.
The essentials of Raft:
- Nodes are leader, follower or candidate. One leader accepts all writes and replicates them to followers as a log.
- A write is committed once a majority of nodes have stored it.
- If the leader stops sending heartbeats, a follower times out, becomes a candidate and asks for votes. A candidate who wins a majority becomes leader for a new term.
- Because any two majorities overlap, there is never more than one leader in a term, and a committed entry survives into the next leader.
Majorities and fault tolerance. A cluster of nodes tolerates failures. Three nodes tolerate one failure, five tolerate two. This is why consensus clusters have an odd number of nodes: a fourth node adds cost without tolerating more failures than three.
You will rarely implement consensus. You will use a coordination service built on it for locks, leader election and configuration, and you should say so rather than reinventing it.
<!--fig:raft-->8. Distributed locks and fencing
A lock service lets one worker at a time do a task. The danger: a worker takes a lock, pauses (a long garbage collection, a network stall), the lock lease expires, another worker takes it, and then the first wakes up and acts, believing it still holds the lock.
The fix is a fencing token: the lock service issues an increasing number with each grant, and the resource being protected rejects any request carrying a lower number than one it has already seen. Without that, a time-based lock does not guarantee exclusion.
9. Transactions across services
A single database gives you transactions. Once data is split across services or shards, you have options, none free:
- Two-phase commit (2PC). A coordinator asks all participants to prepare, then tells them to commit if all agree. It gives atomicity, but participants hold locks while waiting, and a coordinator failure can leave them blocked.
- Sagas. Break the work into local transactions, each with a compensating action that undoes it. If step three fails, run compensations for steps two and one. Sagas avoid long-held locks and give eventual atomicity, but intermediate states are visible and compensations must themselves be idempotent and reliable.
- Outbox pattern. To change a database and publish an event atomically, write the event to an "outbox" table in the same local transaction, and have a separate process publish rows from it. This removes the dual-write problem, where a crash between "update the database" and "send the message" loses one of them.
Prefer designs that keep a business invariant inside one partition so that a local transaction suffices. Choosing the shard key well often removes the need for a distributed transaction.
<!--fig:saga-->Potential deep dives
Deep dive 1: How do you make an operation safe to retry?
The challenge. A call times out. The caller does not know whether it worked, so it retries, and the operation must not happen twice.
Weak: assume the network is reliable. A retried "charge 500" charges twice.
Solid: an idempotency key stored with the result. The client sends a unique key per logical operation. The server records the key and the outcome and returns the stored outcome for a repeat.
Excellent: atomic with the effect, scoped and expiring. Write the key and the effect in one transaction, so a crash cannot leave one without the other. Scope keys to the caller and the operation, reject a reused key with a different request body, and expire old keys after a period longer than any retry window. Where the operation is naturally idempotent ("set the status to shipped"), prefer that form. Pass the key to downstream calls so the whole chain deduplicates.
Deep dive 2: How strong a consistency level do you need?
The challenge. Stronger consistency costs latency and availability, and weaker consistency risks wrong behaviour.
Weak: say "I will make it strongly consistent" everywhere. That pays the highest cost on every operation, including the many where nobody would notice staleness.
Solid: choose per feature. Strong for money and inventory. Read-your-writes for a user's own profile. Eventual for like counts and feeds.
Excellent: connect the choice to mechanisms and failure. Name the mechanism for each level: leader reads or quorums for strong, session tokens for read-your-writes, asynchronous replicas for eventual. Explain what a user would observe in the failure case, such as a partition: a strongly consistent store refuses requests on the minority side, while an eventually consistent one accepts them and may later need to reconcile. Choose conflict resolution deliberately: last-write-wins loses data, application-level merge or conflict-free data types keep both.
Deep dive 3: Do you need a distributed lock?
The challenge. Only one worker should perform a task, or only one process should modify a resource.
Weak: a lock with no expiry. If the holder crashes, nobody can ever take the lock again.
Solid: a lock with a lease and an owner token. The lock expires automatically, and only the owner can release it.
Excellent: question the need, then add fencing. Often an atomic conditional write or a unique constraint removes the need for a lock. If a lock is required, add a fencing token: the lock service issues an increasing number with each grant, and the protected resource rejects any request carrying a lower number than one it has seen. That protects against a holder that pauses past its lease and then resumes. Use a consensus-based coordination service for strong guarantees, and say that a simple cache-based lock is suitable for efficiency, not for strict correctness.
Deep dive 4: Distributed transaction, or saga?
The challenge. One business operation changes data in several services.
Weak: two-phase commit across everything. It holds locks while waiting for every participant, blocks if the coordinator fails, and external services cannot join it.
Solid: a saga with compensations. Local transactions in sequence, each with an undo action, run in reverse when a later step fails.
Excellent: design for visibility and the pivot. Make the saga's state durable and visible, make every step and compensation idempotent and retried, order steps so the hardest to reverse comes last, identify the pivot after which the process must complete, and treat intermediate states as real states in the model. Use the outbox to publish events reliably. Keep invariants inside one partition where possible, so the main flow needs only a local transaction.
What is expected at each level
Mid-level. You know ACID, understand that retries need idempotency, and can describe eventual consistency and quorum reads and writes.
Senior. You choose consistency per feature, design idempotency end to end, explain Raft's majority commit and why consensus clusters are odd-sized, and use sagas and the outbox for cross-service work.
Staff. You discuss which guarantees users and the business actually need, the cost of each, how you would test them under failure, and how to simplify the design so that fewer guarantees are needed at all.
Interview questions and model answers
Q: How do you prevent a customer being charged twice when a client retries? The client sends an idempotency key. The server stores the key and the payment result in one transaction, and returns the stored result on a repeat. Delivery is at-least-once, processing is idempotent, so the effect happens once.
Q: What is the difference between consistency in ACID and in CAP? In ACID it means a transaction moves the database from one valid state to another, respecting invariants. In CAP it means linearizability: every read reflects the most recent write. They are different ideas that share a word.
Q: Why do consensus clusters use an odd number of nodes? A majority is needed to make progress. Three nodes tolerate one failure and four also tolerate only one, so the fourth adds cost but no tolerance.
Q: A job is run by two workers at once even though you use a lock. How? Probably a lease expired during a pause, and the first worker resumed. A fencing token checked by the protected resource prevents the stale worker's writes.
Q: How would you do a multi-service order flow? A saga: reserve inventory, charge payment, create shipment, each as a local transaction with a compensating action. Events are published through an outbox so that state change and message cannot diverge.
Common mistakes
- Claiming exactly-once delivery without saying idempotency makes it so.
- Using wall-clock timestamps to order events across machines.
- Describing a lock as safe without a fencing token.
- Reaching for 2PC across many services by default.
- Choosing an even-sized consensus group.