Consensus Systems in Production: etcd, Consul, and ZooKeeper

LESSON

Consensus and Coordination

014 30 min intermediate

Consensus Systems in Production: etcd, Consul, and ZooKeeper

By the end of this lesson, you will be able to...

  • Turn a control-plane requirement into the watch, lease, session, and revision semantics it needs.

  • Choose a coordination-system shape from explicit workload and failure constraints rather than brand familiarity.

  • Review a client design for stale ownership, missed notifications, and misuse of consensus-backed storage.

Idea in one sentence: Choose etcd, Consul, or ZooKeeper by the authority and client contract you need—revisioned control-plane state, service discovery and health, or session-oriented coordination—not because all three can store data.

Core Insight

The Atlas platform runs many short-lived worker services. Its controllers must agree on the active configuration revision. Services must find healthy peers. One older job framework elects a leader through a session-owned node.

At first, the team sees one storage problem: “we need a distributed key-value store.” That sounds reasonable while the only task is reading a configuration value. It breaks when a controller disconnects during changes, a lease expires while its owner is paused, or a watch misses historical events after compaction. Storage alone does not tell the client what it may safely do.

The real question is not which product is best. It is which coordination contract Atlas needs for each kind of state:

authoritative desired state  -> revision, conditional write, recoverable watch
service reachability         -> registration, health, discovery interface
session-owned coordination   -> session expiry, ephemeral ownership, resync

etcd, Consul, and ZooKeeper can all help build coordination systems, but their natural primitives lead clients toward different designs. The product name is the last part of the decision.

The Promise We Need to Keep

Atlas needs to keep three promises separate:

  1. Configuration authority: controllers must apply changes from one ordered, authoritative version of desired state.
  2. Service discovery: callers should find currently healthy instances without treating a stale membership hint as ownership proof.
  3. Leader ownership: a worker that lost its right to lead must stop acting, even if it wakes up late.

The naive design puts all three into one generic “coordination database,” lets clients watch keys, and treats a successful lock acquisition as permanent authority. This works while connections stay up and the cluster is quiet. It becomes unsafe when the client has partial knowledge.

For example, a watch notification says that something changed. It does not prove the client observed every transition while disconnected. A lease or session says the coordination service may remove an ownership record after its rules are met. It does not physically stop the old process from sending a request. The downstream system needs a revision or fencing value to reject a stale actor.

Start With the Coordination Shape

Before evaluating products, write the smallest contract for each record:

Record Writer Reader reaction Failure boundary Required evidence
desired/version Controller API Reconcile from a known revision Watch disconnect or compaction Current value plus revision
service/payments/instance-7 Service agent Route only to an eligible instance Failed health check or deregistration Discovery/health semantics
leader/reconciler Candidate controller Perform leader-only work Lease/session loss or process pause Current ownership plus fencing value

This table prevents a common category error: treating all changing metadata as the same kind of truth. Desired state needs an authoritative history. Discovery needs a usable view of available instances. Leader election needs a way to deny stale leaders.

The Initial Design: Watch Events as the Whole State

Suppose a controller reads configuration at revision 840, starts a watch, and updates its local cache whenever an event arrives.

read desired state @ revision 840
start watch after 840
apply each received change

This is a good starting shape. It is efficient when the connection remains active and the watch can deliver the requested history.

Now the controller loses its connection for ten minutes. During that time, Atlas writes revisions 841 through 970 and the coordination store compacts older history. The controller reconnects and asks to resume at 841. It cannot safely assume a new watch alone will rebuild its cache; the historical range may be gone.

etcd documents this explicitly: watches stream changes from a chosen revision, but a watcher that starts from or falls behind a compacted revision is canceled. The correction is not “never use watches.” The correction is to make resynchronization part of the client contract:

1. read a complete current snapshot and its revision R
2. reconcile local state with that snapshot
3. watch from R + 1
4. if the watch ends or reports compaction, repeat from step 1

The snapshot is the authority; events are the efficient way to stay current after that boundary.

The Better Model: Products Expose Different Contracts

etcd: revisioned control-plane state

etcd is a natural fit when the main job is a compact, strongly coordinated key space that custom controllers read, compare, update, and watch. Its API exposes revisions, transactions with comparisons, watches from a revision, and leases with TTLs.

In plain English, a revision lets a controller name the state it observed. A conditional transaction lets it say, “write this only if the state is still the version I checked.” A lease lets the cluster associate keys with liveness evidence; if keepalives stop and the lease expires, the attached keys are deleted.

In Atlas, etcd fits desired/version and a controller ownership record when the controller also sends a fencing value to the thing it controls. The lease can remove the old record. The fencing value lets a database, worker, or API reject work from an earlier owner.

Consul: discovery and health-aware coordination

Consul is a natural fit when the central problem is service registration, health, and discovery across a changing service fleet. Its service catalog lets clients discover and monitor service instances; its sessions can be combined with the KV store for advisory locks and client-side leader election.

This does not make every Consul session a proof that its holder is the only process still running. Consul documents lock acquisition as advisory: clients are expected to obey it. Session invalidation can release or delete associated locks, and TTL expiration is a lower bound for invalidation, not a promise that expiry occurs at one exact instant.

That makes Consul a good fit for Atlas's service-discovery promise and for cooperative coordination among well-behaved clients. For an irreversible leader-only side effect, carry Consul's changing lock index or another monotonic fencing value to the downstream resource. Do not rely on elapsed TTL alone.

ZooKeeper: a session-oriented coordination tree

ZooKeeper is a natural fit when an ecosystem already uses its hierarchical znode model, session semantics, ephemeral nodes, sequential nodes, and coordination recipes. An ephemeral znode exists while the creating session remains active and is removed when the session ends. That directly represents “this registration belongs to this live session.”

Its watch model also teaches an important client lesson. Standard ZooKeeper watches are one-time triggers. After receiving an event, a client must set another watch; changes can occur in the gap. While disconnected, the client receives no watches and should act conservatively until it reconnects and resynchronizes.

For Atlas's older job framework, ZooKeeper can therefore be a good compatibility fit. The framework must still reconstruct state after reconnecting and ensure an old leader cannot perform a protected action after losing its session.

Worked Design Review: Atlas Chooses Boundaries

The team now reviews each promise against the system that owns its critical semantics. The following states and timings are illustrative.

| Requirement | Candidate fit | Client rule | What can still fail | | --- | --- | --- | | Controllers reconcile desired/version | etcd | Read snapshot at revision R; watch from R + 1; resync after compaction or reconnect | A paused controller can act on stale cached work unless writes are conditional or fenced | | Callers locate healthy payment instances | Consul | Query discovery/health data; treat it as routing information, not leader authority | Health and reachability observations can lag; requests still need normal retry and failure handling | | Existing batch jobs use session-owned election | ZooKeeper | Re-read election state and re-arm watches after a session or connection event | A former leader may still run locally; protected targets need to reject stale epochs |

The important design choice is not “run three databases for the same record.” Each record has one authoritative owner. Atlas might use etcd for the desired-state control plane, Consul for service discovery, and ZooKeeper only where the legacy framework already owns a separate election contract. Copying the same leader record among them would create competing authority.

So far, we have moved from “pick a distributed store” to “assign one owner and one recovery rule to each coordination fact.” This matters because consensus makes a history defensible, but clients still need to recover after missed observations and stop when their authority becomes stale.

Trade-offs, Limits, and Signals

The trade-off is that consensus-backed coordination buys a clear answer for small, high-value metadata, but writes pay replication and durability costs and clients inherit recovery work.

The boundary appears in a simple question: “If this client wakes up after a partition, what evidence lets it safely resume?” If the answer is only “its old local cache” or “its lease probably has not expired,” the design needs a snapshot, a current revision, or fencing before it acts.

Check Your Understanding

Check: An etcd controller's watch is canceled because its requested revision was compacted. Should it start a new watch at the latest revision and keep its old cache?

Think first, then reveal.

Answer: No. It must first rebuild from a complete current read and record that snapshot's revision. A new watch then keeps that known state current. The old cache may omit changes that were compacted while the controller was away.

Check: A Consul session lock is released after its holder's TTL expires. Can the old holder safely keep issuing leader-only writes?

Think first, then reveal.

Answer: No. The holder may be paused or partitioned and unaware of the expiry. The protected target should verify a current fencing or sequencer value, not trust the former holder's local clock.

Practice: Review a “Universal Coordination Store” Proposal

A team proposes one consensus-backed store for 5 MB customer profiles, clickstream events, a service catalog, a controller's desired state, and a distributed lock. Clients maintain local caches solely from watch events and write leader-only jobs without a fencing value.

Redesign the boundary. Decide which records belong in the coordination system, give the watch-recovery rule, and name the smallest protection needed for leader-only work.

Model answer: Keep small, high-value coordination facts such as desired state, ownership metadata, and perhaps service-catalog metadata in a system whose primitives match their contract. Move customer profiles and clickstream events to stores designed for application data and event throughput. For every cache, read a full snapshot at revision R, watch from R + 1, and resync after disconnect or compaction. Give each successful leader acquisition a monotonic fencing/epoch value and require the protected resource to reject requests with an older value. A lock record alone cannot stop a stale process.

Connections

The previous lesson showed how idempotency and durable operation identity make retries safe. Coordination primitives can decide who may attempt a step, but they do not make the external side effect idempotent; that boundary still needs its own contract.

The next lesson uses Jepsen-style verification to turn these API claims into histories, faults, and invariants. A watch resync rule, lease-expiry response, and fencing check are all testable client-visible contracts.

Resources

Key Takeaways

PREVIOUS Exactly-Once Semantics, Idempotency, and Deduplication NEXT Jepsen-Style Verification and Failure Injection