Consensus Systems in Production: etcd, Consul, and ZooKeeper
LESSON
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:
- Configuration authority: controllers must apply changes from one ordered, authoritative version of desired state.
- Service discovery: callers should find currently healthy instances without treating a stale membership hint as ownership proof.
- 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.
- Watch-based designs cost resync logic. A reconnect is not proof that no update was missed. Track watch cancellations, compaction responses, and time since the last successful snapshot.
- Lease and session designs cost careful authority checks. They improve cleanup of ownership records, but a paused or partitioned process can continue locally. Track expired sessions, lease keepalive failures, rejected fencing values, and leader changes.
- Discovery optimizes routing, not arbitrary transactions. Treat health data as a decision input under its documented freshness model, not as a substitute for an application's own correctness checks.
- Small authoritative records are the intended workload. Large documents, high-volume telemetry, queue payloads, and hot business data make quorum latency and storage pressure part of the application hot path.
- Multiple systems cost operational clarity. Use more than one only when their authority boundaries are explicit; otherwise the integration itself becomes a coordination problem.
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
- [DOC] etcd API — Focus: Revisions, transactions, watches, compaction responses, and TTL leases.
- [DOC] Consul: Service Discovery — Focus: Service catalog, health, and discovery semantics.
- [DOC] Consul Sessions and Distributed Locks — Focus: Advisory locks, session invalidation, TTL limits, and the lock index as a sequencer.
- [DOC] ZooKeeper Programmer's Guide — Focus: Session-bound ephemeral nodes, one-time watches, and safe behavior across disconnection.
Key Takeaways
- etcd, Consul, and ZooKeeper are coordination systems whose API semantics matter more than their shared ability to store records.
- A watch is an incremental view after a known snapshot; reconnect and compaction require resynchronization.
- Leases and sessions clean up ownership records, but fencing is what lets a protected resource reject a stale actor.
- Use consensus-backed state for small facts whose disagreement is dangerous, and give every fact one authoritative owner.
- A product choice is defensible only when its revisions, health model, watches, sessions, and recovery behavior match the workload's contract.