Designing Data-Intensive Applications
Ch. 9

Distributed Transactions and Consensus

Consensus protocols let nodes agree on a value despite failures — powering distributed transactions.

Getting multiple nodes to agree on a single value — who is leader, whether to commit, what the next log entry is — is the consensus problem. Algorithms like Raft make this practical for production systems.

Diagram
In practice

etcd and HashiCorp Consul use Raft for leader election and consistent key-value storage — the backbone of Kubernetes control planes. Kafka replaced ZooKeeper with KRaft (internal Raft) for broker metadata. Cross-database 2PC (XA transactions) exists in Java EE but is fragile; the saga pattern with compensating transactions in microservices is more common today.

Google at scale

Chubby (precursor to etcd) provided distributed consensus for Google's internal infrastructure — leader election, lock service, and membership. Kubernetes today relies on etcd's Raft implementation for the same guarantees.

typescript — SAGA compensating transaction
// Microservice SAGA — compensate if payment fails (Airbnb-style)
async function bookListing(input: BookingInput) {
  const hold = await inventoryService.reserve(input); // step 1
  try {
    const payment = await paymentService.charge(input); // step 2
    await notificationService.sendConfirmation(input); // step 3
    return { hold, payment };
  } catch (err) {
    await inventoryService.release(hold.id); // compensating transaction
    throw err;
  }
}
Key Takeaways
  • Two-phase commit (2PC) coordinates atomic commit across participants.
  • 2PC blocks if the coordinator fails after prepare.
  • Fault-tolerant consensus (Raft, Paxos, Zab) elects a leader and replicates a log.
  • Consensus is required for leader election, atomic commit, and membership changes.
  • Coordination services (ZooKeeper, etcd) provide consensus as a service.
  • Raft in etcd/Consul, Kafka KRaft, and XA transactions across PostgreSQL show consensus in production.
consensusRaftetcdKafkatwo-phase commitZooKeeper