All three are small, replicated, strongly consistent stores for a few megabytes of critical metadata, so they get compared as interchangeable "coordination services". They are not interchangeable in what they are for. etcd is a consistent key-value store with watches and leases, built as the substrate for control planes like Kubernetes. ZooKeeper is the original coordination kernel, a tree of znodes with sessions and ephemeral nodes that a generation of data systems embedded. Consul is a service networking product — service catalog, health checks, DNS, a KV store and a mesh — with consensus inside it and gossip around it. The question that decides it is: are you building a control plane that needs a consistent store, or do you need service discovery as a product? A consistent store points to etcd (or the one your platform already runs). Discovery and health checking across VMs and datacenters points to Consul. ZooKeeper is the answer mainly when the software you are running already requires it.
The Verdict#
Default: etcd for new coordination needs, and only if a conditional write in your existing database will not do. Consul when the requirement is service discovery and health checking outside Kubernetes. ZooKeeper when a system you depend on still requires it.
| Pick etcd when | Pick ZooKeeper when | Pick Consul when |
|---|---|---|
| You are building a control plane: leader election, shard maps, config with watches | You run software that depends on it (older Kafka, HBase, Solr, some ClickHouse deployments) | You need a service catalog, health checks and DNS across VMs, containers and datacenters |
| You are on Kubernetes and want the same operational model as its control plane | You have deep in-house ZooKeeper expertise and recipes built on ephemeral sequential znodes | You want discovery that keeps answering during partitions (stale reads by default for DNS) |
| You want linearizable reads by default and a simple gRPC API | You need hierarchical namespaces with per-node ACLs and a mature Java client ecosystem | You want KV, locks, discovery and mesh from one product with one operational model |
| Data stays under ~2 GiB (default quota), well under the 8 GiB suggested maximum | Your workload is read-heavy and can tolerate reads served by any server unless you sync | Multi-datacenter federation is a first-class requirement |
🎯 Staff Move: "For leader election and the shard map I'd use etcd with a 10-second lease and fencing tokens — or, honestly, a conditional write in Postgres, since we already run it. Service discovery is a different problem: it wants availability during partitions, so I'd use the platform's discovery, not a hand-rolled registry in a consensus store."
When the non-default wins:
- ZooKeeper for new work wins only when you are extending a system that already runs it well and has client recipes battle-tested against it; adding a second consensus system to avoid it is worse.
- Consul for config and locks wins when Consul is already deployed for discovery on every node: one fewer system beats a slightly better fit.
- None of them wins more often than people expect. Leader election for a handful of singleton jobs works with a conditional update on a lease row in the database you already run, with the row's epoch as the fencing token.
At a Glance#
| Dimension | etcd | ZooKeeper | Consul |
|---|---|---|---|
| Data model | Flat key space with prefixes; MVCC revisions per key | Hierarchical tree of znodes; persistent, ephemeral, sequential | Service catalog + health checks + flat KV (prefixes) |
| Consensus | Raft | ZAB (ZooKeeper Atomic Broadcast) | Raft among 3–5 servers; gossip (Serf/SWIM) among all agents |
| Read consistency | Linearizable by default; serializable (local) reads optional | Reads served by the connected server, may be stale; sync before read for freshness | default (leader-served, small stale window on leader change), consistent, or stale; DNS uses stale |
| Ordering | Global revision number on every write | Global zxid on every write; FIFO per client session | Raft index per write |
| Liveness primitive | Leases with TTL, kept alive by the client | Sessions; ephemeral znodes deleted when session expires | Agent health checks + gossip failure detection; sessions for locks |
| Change notification | Watches on keys or prefixes, resumable from a revision | Watches (one-shot classic; persistent and recursive since 3.6) | Blocking queries (long poll on an index); watches built on them |
| Write throughput | ~10K–50K small writes/s on good SSDs (bounded by leader fsync) | ~10K–20K writes/s; reads scale with servers | ~thousands of KV writes/s; catalog churn drives load |
| Write latency | ~2–10ms in one region (quorum fsync) | ~2–10ms in one region | ~5–20ms; consistent reads add a quorum round trip |
| Size limits (as of 2026) | 1.5 MiB per request default; 2 GiB storage default, 8 GiB suggested max | ~1 MB per znode default (jute.maxbuffer); whole dataset in heap | 512 KB per KV value |
| Scaling model | 3 or 5 voters; learners for catch-up; split keyspaces into separate clusters | 3 or 5 voters + observers for read and fan-out scaling | 3 or 5 servers per datacenter; thousands of client agents; federate datacenters |
| Operational burden | Medium: disk latency, compaction and defrag, quota | Medium-high: JVM tuning and GC pauses, session timeouts, four-letter-word monitoring | Medium-high: servers, agents everywhere, ACLs, gossip encryption, upgrades |
| Managed options | Embedded in every managed Kubernetes control plane; few standalone services | Bundled in some managed data services; rarely standalone | Vendor-managed service; self-host widely |
| Cost shape | 3–5 small VMs with fast SSDs; cost is people | Same, plus JVM memory | Servers plus an agent per node; enterprise license for some features |
| License | Apache 2.0 (CNCF) | Apache 2.0 (ASF) | Business Source License since 2023 |
Numbers to bring:
| Figure | Value | Condition |
|---|---|---|
| Voters | 3 tolerate 1 failure, 5 tolerate 2 | Never even numbers |
| etcd storage quota | 2 GiB default, 8 GiB suggested max | As of 2026 |
| etcd request size | 1.5 MiB default max | As of 2026 |
| ZooKeeper znode size | ~1 MB default limit | Configurable, but keep values in KB |
| Consul KV value | 512 KB max | As of 2026 |
| Leader-election lease TTL | 10–15s, keepalive every ~3s | Typical choice |
| Consensus write latency | ~2–10ms | One region, SSDs |
| etcd fsync alert | WAL fsync p99 above ~10ms | Disk is the usual root cause of elections |
How They Actually Differ#
What all three share, so you can skip it in the interview:
- A leader-based consensus protocol with majority quorums: writes go through one leader and are fsynced on a majority before acknowledgment.
- The whole dataset is meant to be small and largely memory-resident; they are metadata stores, not databases.
- Liveness by heartbeat (lease, session or health check), which means pauses longer than the timeout look like death.
- Watches or blocking queries instead of polling, and a monotonically increasing version number usable as a fencing token.
1. Store vs Product#
etcd and ZooKeeper are building blocks: they give you primitives — compare-and-swap, leases or sessions, watches, ordered writes — and you build leader election, locks or membership on top, ideally with a well-tested client library recipe. Consul is a product with an opinion: a node runs an agent, services register with health checks, and other services find healthy instances through DNS or HTTP. The KV store and locks are there, but they are secondary.
That changes who owns what. With etcd you own the recipes and the correctness of every client. With Consul you own an agent on every node and a product's upgrade cadence, but discovery and health checking work without writing code.
2. Reads: Who Answers and How Fresh#
| System | Default read path | Freshness | Cost |
|---|---|---|---|
| etcd | Linearizable: leader confirms it is still leader (ReadIndex) before answering | Always current | One extra network round to a quorum |
| ZooKeeper | The server the client is connected to answers locally | May lag the leader; ordered per client | Cheapest; scales with servers and observers |
| Consul HTTP | Leader answers using its leader lease | Current except a brief window during leader change | One hop to the leader |
| Consul DNS | Any server answers (stale mode) | Can lag by the replication delay | Scales with servers; survives losing the leader |
This is the most interview-relevant difference. ZooKeeper's local reads make it fast and read-scalable but surprise people who assumed linearizability: a client can read a value that another client has already seen overwritten, unless it calls sync first. etcd made the opposite default. Consul chose a split default that matches its purpose: catalog lookups through DNS favor availability, and API reads favor consistency.
🎯 Staff Insight: Service discovery is an availability problem wearing a consistency costume. A slightly stale list of healthy instances is fine; no list at all during a leader election is an outage. That is why Consul's DNS reads are stale by default — and why putting discovery on linearizable reads is the wrong instinct.
3. Liveness: Leases, Sessions and Gossip#
- etcd leases: the client creates a lease with a TTL (commonly 10–15s for leader election) and keeps it alive; keys attached to the lease vanish when it expires. One lease can cover many keys, so 1,000 workers cost 1,000 keepalives, not one per key.
- ZooKeeper sessions: the client holds a session with a negotiated timeout (commonly 6–30s); ephemeral znodes die with the session. A long JVM GC pause on the client expires its session and drops its locks and registrations — the classic ZooKeeper production incident.
- Consul: liveness of services comes from health checks run by the local agent (HTTP, TCP, script, TTL), plus gossip-based failure detection between agents with suspicion and refutation to avoid flapping. This scales to thousands of nodes without each node heartbeating the servers, which is the main reason Consul handles large fleets better than a registry built on raw sessions or leases.
The fencing step is identical for all three systems: the lock service cannot stop a paused ex-leader, only the resource can.
4. Scale and Topology#
All three have the same consensus arithmetic: 3 voters tolerate 1 failure, 5 tolerate 2, even numbers add latency without adding tolerance, and adding voters slows writes. They differ in how they scale everything else.
| Need | etcd | ZooKeeper | Consul |
|---|---|---|---|
| More read capacity | Serializable reads from followers, or a cache in front (Kubernetes API server's watch cache) | Observers: non-voting servers that serve reads and watches | More servers in stale mode; agents cache |
| Many clients | Each client holds a gRPC connection; watch fan-out costs leader CPU | Each client holds a session; observers absorb fan-out | Agents per node; gossip spreads membership; servers see agents, not apps |
| Multiple regions | One cluster per region; cross-region voters make every write a WAN round trip | Same; observers in remote regions for reads | Federated datacenters, each with its own servers; WAN gossip between them |
| More data | Split by function into separate clusters | Same; also heap-bounded | Move non-config data elsewhere |
Consul's multi-datacenter model is the clearest built-in story of the three: each datacenter has its own Raft cluster, so losing a WAN link does not stop local writes, and queries can be forwarded across datacenters when needed.
Where Each One Breaks#
| System | Failure mode | Symptom | Detection | Mitigation | Owner |
|---|---|---|---|---|---|
| etcd | NOSPACE alarm | Cluster goes read-only for writes; control plane stops changing | etcd_mvcc_db_total_size_in_bytes vs quota | Auto-compaction, scheduled defrag, raise quota toward 8 GiB, stop storing bulk data | Platform |
| etcd | Slow disk → leader elections | Leader changes every few minutes; write p99 from 5ms to 500ms | etcd_disk_wal_fsync_duration_seconds p99 over ~10ms, leader change count | Dedicated SSDs, no noisy neighbors | Platform |
| etcd | Watch fan-out overload | Thousands of clients watching large prefixes; leader CPU saturates | Watcher count, leader CPU | Put a cache layer in front; narrow prefixes | Platform + clients |
| ZooKeeper | Session expiry storm | A network blip or GC pause expires hundreds of sessions; ephemeral nodes vanish; mass rebalance | Session expirations per minute; zk_outstanding_requests | Longer session timeouts for non-critical clients, GC tuning, jittered reconnect | Clients + platform |
| ZooKeeper | Server GC pause | Server stalls, followers fall behind, leader election | JVM GC pause time, zk_avg_latency | Heap sizing, low-pause collector, dedicated hosts | Platform |
| ZooKeeper | Stale read assumption | Client reads old config after another client saw the new one | Hard to detect; correctness bug | sync before critical reads; version checks on writes | Client teams |
| Consul | Central cluster overload | Servers saturated by catalog churn and KV misuse; everything that depends on discovery fails | Raft commit time, server CPU, RPC rate | Separate clusters per function, move non-config KV out, rate-limit clients | Platform |
| Consul | Storage engine pathology | Raft log writes slow down as freelist bookkeeping grows; leader cannot keep up | Raft commit latency, disk write amplification | Keep storage engine current; snapshot and trim logs; capacity reviews | Platform |
| Consul | Circular dependencies | Telemetry and deploy tooling depend on Consul; during the outage you are blind and cannot deploy the fix | Game days | Break-glass paths that do not depend on discovery | Principal-level |
| All three | Quorum loss | Two of three voters down: no writes, no leader election | Has-leader metric | 5 voters across 3 zones for tier-0; tested restore from snapshot | Platform |
The production surprise for each:
- etcd: nobody watches the database size until the
NOSPACEalarm stops the control plane on a Friday. - ZooKeeper: the incident is usually caused by a client GC pause, not the ensemble.
- Consul: it becomes the dependency of everything — discovery, secrets bootstrap, config, deploys — so its outage is a company outage. Blast radius is the design problem.
Incident Sketch: The Session Expiry Storm#
t=0 Top-of-rack switch flaps for 9 seconds in one zone
t=+6s ZooKeeper sessions with a 6s timeout expire for 400 workers in that zone
t=+7s Their ephemeral znodes vanish; the coordinator sees 400 departures
t=+8s Coordinator reassigns 3,000 partitions to the remaining zones
t=+12s Network recovers; 400 workers reconnect with new sessions and rejoin
t=+13s Coordinator reassigns again; two full rebalances in 10 seconds
t=+5min Consumers still catching up; downstream lag alarms
Detection: session expirations per minute, rebalance count, membership churn. Mitigation: longer session timeouts for members whose loss is cheap to delay, a grace period before reassigning work, and rebalance rate limiting. Owner: the client team that wrote the membership recipe, with platform setting the defaults.
Incident Sketch: etcd Hits Its Quota#
t=0 A new controller writes a status key per object every 10 seconds
t=+5d etcd database grows from 600 MB to 2 GB; compaction keeps revisions, defrag never scheduled
t=+5d+1m NOSPACE alarm: etcd rejects writes
t=+5d+2m Control plane cannot create, update or delete objects; running workloads continue
t=+5d+40m Compact, defragment member by member, disarm alarm; offending controller rate-limited
Detection that would have caught it: database size as a percentage of quota with an alert at 60%, writes per second by client, and keys per prefix. Prevention: per-client write budgets, scheduled defragmentation, and a review for any controller that writes status at a fixed interval.
Cost and Operations#
| etcd | ZooKeeper | Consul | |
|---|---|---|---|
| Who runs it | Platform team; or the managed Kubernetes provider for its own control plane | Platform team, often the team that runs the dependent data system | Platform / networking team; agents on every node owned by the fleet team |
| What the bill scales with | Almost nothing: 3–5 small VMs with fast disks; people and on-call | Same, plus JVM memory and tuning time | Servers per datacenter plus agent footprint per node; license for enterprise features |
| Typical footprint | 3–5 nodes, 2–8 vCPU, 8–32 GB RAM, low-latency SSD | 3–5 nodes, similar; heap a few GB | 3–5 servers per datacenter; an agent on each of thousands of nodes |
| Hidden cost | Every team building its own recipes; consolidating clusters later | Expertise is getting rarer as dependent systems remove it | Correlated failure: one product underneath discovery, config and mesh |
The infrastructure bill is trivial for all three — a few hundred to a few thousand dollars a month. The cost is on-call and expertise. A consensus cluster is a tier-0 dependency that needs someone who understands quorum loss at 3 a.m. Running two different ones doubles that.
| Scale | Sensible setup | Infrastructure | People |
|---|---|---|---|
| Startup (1 region, a few singleton jobs) | Lease row in the main database; Kubernetes-provided etcd for the cluster | ~$0 extra | None dedicated |
| Growth (several control planes, VMs plus Kubernetes) | One shared etcd for platform control planes; Consul for VM discovery | Hundreds of dollars a month | 1–2 platform engineers part-time |
| Large (multi-region, thousands of nodes) | Cells: separate clusters per region and per function; one standard library for election and fencing | Low thousands a month | A platform team that owns coordination as a product |
Assumptions: small VMs with SSDs at cloud list prices; the dominant cost at every tier is expertise and on-call.
🧭 Principal Insight: The org-level question is how many consensus systems you run, not which one is best. Each extra one is another tier-0 dependency with its own failure modes, upgrade path and expertise. Standardize on one per use case and retire the rest.
Switching Later#
| Migration | Difficulty | What's hard to undo |
|---|---|---|
| ZooKeeper → etcd (own code) | Moderate: recipes map (ephemeral node → leased key, sequential node → revision), but every client changes | Assumptions about local reads, one-shot watches and session semantics |
| ZooKeeper → built-in consensus (vendor product) | Follows the vendor: e.g. Kafka's move to its internal Raft quorum, ClickHouse Keeper's ZooKeeper-compatible protocol | Usually a one-way migration with a documented procedure; plan a maintenance window |
| Consul → Kubernetes-native discovery | Moderate to hard: DNS names, health checks and mesh policies everywhere | Workloads outside Kubernetes (VMs, databases) lose discovery |
| etcd/ZooKeeper registry → Consul | Moderate | Agent rollout to every node; new operational model |
| Any → database-as-coordinator (conditional writes) | Easy for leader election and locks | Little; often the simplest long-term answer for low-rate coordination |
Retiring a ZooKeeper-based recipe, in order:
- Put the recipe behind a small internal interface (
elect(),heartbeat(),currentLeader()) in every client. - Implement the interface on the new backend with the same fencing token semantics.
- Run both for one role at a time, with the new backend authoritative only after a clean week.
- Remove the old backend per role, then decommission the ensemble once no client connects.
The one-way doors:
- Agents on every node (Consul). Once discovery, health checks and mesh depend on them, removing the agent is a fleet-wide migration.
- Hand-written recipes against a specific API. Leader election code that depends on ZooKeeper's sequential ephemeral znodes does not port; code that depends on a small internal interface does.
- Putting application data in the coordination store. Getting it out later means finding every reader, and the store's size problems arrive first.
How Real Companies Chose#
Kubernetes — etcd as the Single Backing Store#
Kubernetes keeps all cluster data in etcd, which its documentation describes as a consistent and highly available key-value store and the backing store for all cluster data. It recommends an odd number of members, five for production, dedicated machines or isolated environments, and a backup plan, and warns that disk and network starvation cause heartbeat timeouts and instability — without a leader, no cluster state can change (Kubernetes docs).
Staff insight: etcd's flagship user runs it as a dedicated, five-member, backed-up tier-0 store behind a caching API server. That is the template for any control plane you design.
Apache Kafka — Replacing ZooKeeper with a Built-In Quorum#
KIP-500 replaced Kafka's ZooKeeper dependency with a self-managed, Raft-based metadata quorum. The stated reasons: operators had to learn and secure two distributed systems; metadata held outside Kafka could diverge from controller state; and controller failover required reloading all state, while a quorum of controllers that already track the metadata log can fail over quickly and support far more partitions (KIP-500).
Staff insight: For a product, an external coordination service is a tax on every operator. If you are building the product, embed consensus; if you are buying products, prefer the ones that already did. ClickHouse made a similar move with ClickHouse Keeper, a Raft-based, ZooKeeper-protocol-compatible replacement with linearizable reads (ClickHouse docs).
Roblox — When Consul Became a Single Point of Failure#
Roblox's 73-hour outage in October 2021 centered on Consul, which served service discovery for thousands of services and underpinned Nomad and Vault. Two causes combined: a newly enabled streaming feature created contention under heavy concurrent read and write load, and a BoltDB freelist pathology made each Raft log append rewrite megabytes of bookkeeping. Afterward Roblox split critical workloads into dedicated Consul clusters, moved inappropriate data out of the KV store, removed circular dependencies between telemetry and Consul, and moved to bbolt (Roblox postmortem).
Staff insight: The product choice was not the root problem; the blast radius was. One coordination cluster underneath discovery, scheduling, secrets and telemetry turns any bug into a total outage. Cells, separate clusters per function, and break-glass paths are the design.
Follow-Ups to Expect#
| After You Say... | They Will Ask... | What They're Testing |
|---|---|---|
| "etcd for leader election" | "The leader pauses for 20 seconds. What happens?" | Leases expire; fencing tokens checked by the resource |
| "etcd" | "Why not just use a conditional write in our database?" | Whether you add a tier-0 dependency only when it pays |
| "ZooKeeper" | "Are reads linearizable?" | Local reads, sync, and when staleness matters |
| "Consul for discovery" | "What happens to discovery if Consul loses quorum?" | Stale reads, agent caches, clients keeping last-known endpoints |
| "Consul" | "What else depends on it, and what's the blast radius?" | Correlated failure, separate clusters, break-glass paths |
| "5 voters" | "Why not 6? Why not 7 across 7 regions?" | Quorum arithmetic and WAN write latency |
| "Watches for config" | "10,000 clients watch the same prefix. Then what?" | Fan-out cost, caching layer, observers |
| "Store the shard map in etcd" | "How big can it get?" | 2 GiB default quota, 8 GiB suggested max, compaction and defrag |
| "Consul agents on every node" | "An agent's host is overloaded. What does gossip do?" | Failure detection with suspicion and refutation; false positives under CPU starvation |
| "Leases with a 10s TTL" | "Why not 1 second?" | Pause tolerance vs failover time; GC and network jitter |
| "Kubernetes discovery" | "What about the databases running on VMs?" | Hybrid environments and where Kubernetes-native discovery stops |
What to Say in the Interview#
"These are metadata stores, not databases. I'll keep data in them to kilobytes per key and a few hundred megabytes total, and keep them off the per-request path so losing quorum freezes changes but not traffic."
"For a new control plane I'd pick etcd because reads are linearizable by default and it's what Kubernetes already runs. For discovery across VMs and datacenters I'd pick Consul, because discovery needs to keep answering during a partition."
"Every lock is a lease, and every lease needs a fencing token checked by the resource — that's true whichever of the three we choose."
"The bigger risk is blast radius: if discovery, secrets and deploys all sit on one cluster, its bad day is our outage. I'd separate clusters by function and keep a break-glass path."
Related Guides#
- etcd and ZooKeeper — internals, recipes, fencing and failure modes in depth
- Design a Consensus Service — building the store itself
- Design a Service Registry — discovery as an availability problem
- Design Distributed Locking — leases, fencing and when not to lock
- Design a Job Scheduler — singleton workers and leader election in practice
- Kubernetes — the canonical control-plane / data-plane split on etcd
- Coordination Strategies — choosing between consensus, conditional writes and avoidance
- Feature Flags — config distribution with watches and staged rollout