Hiring BarSupport

etcd vs ZooKeeper vs Consul

Comparison19 min read4 diagrams

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 whenPick ZooKeeper whenPick Consul when
You are building a control plane: leader election, shard maps, config with watchesYou 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 planeYou have deep in-house ZooKeeper expertise and recipes built on ephemeral sequential znodesYou want discovery that keeps answering during partitions (stale reads by default for DNS)
You want linearizable reads by default and a simple gRPC APIYou need hierarchical namespaces with per-node ACLs and a mature Java client ecosystemYou 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 maximumYour workload is read-heavy and can tolerate reads served by any server unless you syncMulti-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.
Diagram: The Verdict

At a Glance#

DimensionetcdZooKeeperConsul
Data modelFlat key space with prefixes; MVCC revisions per keyHierarchical tree of znodes; persistent, ephemeral, sequentialService catalog + health checks + flat KV (prefixes)
ConsensusRaftZAB (ZooKeeper Atomic Broadcast)Raft among 3–5 servers; gossip (Serf/SWIM) among all agents
Read consistencyLinearizable by default; serializable (local) reads optionalReads served by the connected server, may be stale; sync before read for freshnessdefault (leader-served, small stale window on leader change), consistent, or stale; DNS uses stale
OrderingGlobal revision number on every writeGlobal zxid on every write; FIFO per client sessionRaft index per write
Liveness primitiveLeases with TTL, kept alive by the clientSessions; ephemeral znodes deleted when session expiresAgent health checks + gossip failure detection; sessions for locks
Change notificationWatches on keys or prefixes, resumable from a revisionWatches (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 heap512 KB per KV value
Scaling model3 or 5 voters; learners for catch-up; split keyspaces into separate clusters3 or 5 voters + observers for read and fan-out scaling3 or 5 servers per datacenter; thousands of client agents; federate datacenters
Operational burdenMedium: disk latency, compaction and defrag, quotaMedium-high: JVM tuning and GC pauses, session timeouts, four-letter-word monitoringMedium-high: servers, agents everywhere, ACLs, gossip encryption, upgrades
Managed optionsEmbedded in every managed Kubernetes control plane; few standalone servicesBundled in some managed data services; rarely standaloneVendor-managed service; self-host widely
Cost shape3–5 small VMs with fast SSDs; cost is peopleSame, plus JVM memoryServers plus an agent per node; enterprise license for some features
LicenseApache 2.0 (CNCF)Apache 2.0 (ASF)Business Source License since 2023

Numbers to bring:

FigureValueCondition
Voters3 tolerate 1 failure, 5 tolerate 2Never even numbers
etcd storage quota2 GiB default, 8 GiB suggested maxAs of 2026
etcd request size1.5 MiB default maxAs of 2026
ZooKeeper znode size~1 MB default limitConfigurable, but keep values in KB
Consul KV value512 KB maxAs of 2026
Leader-election lease TTL10–15s, keepalive every ~3sTypical choice
Consensus write latency~2–10msOne region, SSDs
etcd fsync alertWAL fsync p99 above ~10msDisk 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.

Diagram: 1. Store vs Product

2. Reads: Who Answers and How Fresh#

SystemDefault read pathFreshnessCost
etcdLinearizable: leader confirms it is still leader (ReadIndex) before answeringAlways currentOne extra network round to a quorum
ZooKeeperThe server the client is connected to answers locallyMay lag the leader; ordered per clientCheapest; scales with servers and observers
Consul HTTPLeader answers using its leader leaseCurrent except a brief window during leader changeOne hop to the leader
Consul DNSAny server answers (stale mode)Can lag by the replication delayScales 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.
Diagram: 3. Liveness: Leases, Sessions and Gossip

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.

NeedetcdZooKeeperConsul
More read capacitySerializable reads from followers, or a cache in front (Kubernetes API server's watch cache)Observers: non-voting servers that serve reads and watchesMore servers in stale mode; agents cache
Many clientsEach client holds a gRPC connection; watch fan-out costs leader CPUEach client holds a session; observers absorb fan-outAgents per node; gossip spreads membership; servers see agents, not apps
Multiple regionsOne cluster per region; cross-region voters make every write a WAN round tripSame; observers in remote regions for readsFederated datacenters, each with its own servers; WAN gossip between them
More dataSplit by function into separate clustersSame; also heap-boundedMove 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#

SystemFailure modeSymptomDetectionMitigationOwner
etcdNOSPACE alarmCluster goes read-only for writes; control plane stops changingetcd_mvcc_db_total_size_in_bytes vs quotaAuto-compaction, scheduled defrag, raise quota toward 8 GiB, stop storing bulk dataPlatform
etcdSlow disk → leader electionsLeader changes every few minutes; write p99 from 5ms to 500msetcd_disk_wal_fsync_duration_seconds p99 over ~10ms, leader change countDedicated SSDs, no noisy neighborsPlatform
etcdWatch fan-out overloadThousands of clients watching large prefixes; leader CPU saturatesWatcher count, leader CPUPut a cache layer in front; narrow prefixesPlatform + clients
ZooKeeperSession expiry stormA network blip or GC pause expires hundreds of sessions; ephemeral nodes vanish; mass rebalanceSession expirations per minute; zk_outstanding_requestsLonger session timeouts for non-critical clients, GC tuning, jittered reconnectClients + platform
ZooKeeperServer GC pauseServer stalls, followers fall behind, leader electionJVM GC pause time, zk_avg_latencyHeap sizing, low-pause collector, dedicated hostsPlatform
ZooKeeperStale read assumptionClient reads old config after another client saw the new oneHard to detect; correctness bugsync before critical reads; version checks on writesClient teams
ConsulCentral cluster overloadServers saturated by catalog churn and KV misuse; everything that depends on discovery failsRaft commit time, server CPU, RPC rateSeparate clusters per function, move non-config KV out, rate-limit clientsPlatform
ConsulStorage engine pathologyRaft log writes slow down as freelist bookkeeping grows; leader cannot keep upRaft commit latency, disk write amplificationKeep storage engine current; snapshot and trim logs; capacity reviewsPlatform
ConsulCircular dependenciesTelemetry and deploy tooling depend on Consul; during the outage you are blind and cannot deploy the fixGame daysBreak-glass paths that do not depend on discoveryPrincipal-level
All threeQuorum lossTwo of three voters down: no writes, no leader electionHas-leader metric5 voters across 3 zones for tier-0; tested restore from snapshotPlatform

The production surprise for each:

  • etcd: nobody watches the database size until the NOSPACE alarm 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#

etcdZooKeeperConsul
Who runs itPlatform team; or the managed Kubernetes provider for its own control planePlatform team, often the team that runs the dependent data systemPlatform / networking team; agents on every node owned by the fleet team
What the bill scales withAlmost nothing: 3–5 small VMs with fast disks; people and on-callSame, plus JVM memory and tuning timeServers per datacenter plus agent footprint per node; license for enterprise features
Typical footprint3–5 nodes, 2–8 vCPU, 8–32 GB RAM, low-latency SSD3–5 nodes, similar; heap a few GB3–5 servers per datacenter; an agent on each of thousands of nodes
Hidden costEvery team building its own recipes; consolidating clusters laterExpertise is getting rarer as dependent systems remove itCorrelated 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.

ScaleSensible setupInfrastructurePeople
Startup (1 region, a few singleton jobs)Lease row in the main database; Kubernetes-provided etcd for the cluster~$0 extraNone dedicated
Growth (several control planes, VMs plus Kubernetes)One shared etcd for platform control planes; Consul for VM discoveryHundreds of dollars a month1–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 fencingLow thousands a monthA 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#

MigrationDifficultyWhat's hard to undo
ZooKeeper → etcd (own code)Moderate: recipes map (ephemeral node → leased key, sequential node → revision), but every client changesAssumptions 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 protocolUsually a one-way migration with a documented procedure; plan a maintenance window
Consul → Kubernetes-native discoveryModerate to hard: DNS names, health checks and mesh policies everywhereWorkloads outside Kubernetes (VMs, databases) lose discovery
etcd/ZooKeeper registry → ConsulModerateAgent rollout to every node; new operational model
Any → database-as-coordinator (conditional writes)Easy for leader election and locksLittle; often the simplest long-term answer for low-rate coordination
Diagram: Switching Later

Retiring a ZooKeeper-based recipe, in order:

  1. Put the recipe behind a small internal interface (elect(), heartbeat(), currentLeader()) in every client.
  2. Implement the interface on the new backend with the same fencing token semantics.
  3. Run both for one role at a time, with the new backend authoritative only after a clean week.
  4. Remove the old backend per role, then decommission the ensemble once no client connects.

The one-way doors:

  1. Agents on every node (Consul). Once discovery, health checks and mesh depend on them, removing the agent is a fleet-wide migration.
  2. 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.
  3. 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."

  1. Loading the index…