Designing Data-Intensive Applications
Ch. 8

Knowledge, Truth, and Lies

In a distributed system, nodes operate with incomplete information — and must make decisions anyway.

A node cannot know global system state — it only knows what it has observed. Decisions based on partial knowledge (is the leader dead? did my write succeed?) are inherently uncertain. Protocols define rules for reaching sufficient agreement to act.

In practice

A stale PostgreSQL primary that was network-partitioned must be fenced before rejoining — otherwise it could accept writes that conflict with the new primary. etcd grants expiring leases so only the current leader holds the lock. DynamoDB conditional writes (attribute_not_exists) prevent lost-update races without a central lock service.

Shopify at scale

During a network partition, the old primary must not resume writes after a new leader is elected. Fencing tokens from etcd ensure stale nodes cannot corrupt shared inventory or checkout state.

typescript — Fencing stale leaders
// Shopify leader election — fencing token blocks stale primary writes
const token = await etcd.getLease().then((l) => l.fencingToken);
async function writeToSharedStore(key: string, value: string) {
  const current = await storage.getFencingToken(key);
  if (token < current) throw new Error("stale leader — fenced");
  await storage.put(key, value, { fencingToken: token });
}
Key Takeaways
  • The majority quorum defines what is true in leaderless systems.
  • Byzantine faults: nodes that deliberately lie or behave arbitrarily.
  • System models (crash-stop vs Byzantine) assume different failure modes.
  • Fencing tokens prevent stale leaders from corrupting shared resources.
  • Truth is agreement among nodes — not an absolute property.
  • etcd lease fencing, ZooKeeper ephemeral nodes, and DynamoDB conditional writes enforce safe decisions.
Byzantine faultetcdZooKeeperfencing tokenDynamoDBquorum