10.4 Coordination Services
ZooKeeper, etcd, Consul — modeled after Google's Chubby lock service.
Although they look superficially like any other key-value store, THEY ARE NOT DESIGNED FOR HIGH WRITE VOLUMES OR GENERAL-PURPOSE DATA STORAGE. Instead, they are designed TO COORDINATE AMONG NODES OF ANOTHER DISTRIBUTED SYSTEM. (Kubernetes relies on etcd; Spark and Flink in HA mode rely on ZooKeeper.)
They hold SMALL amounts of data that FIT ENTIRELY IN MEMORY (although they still write to disk for durability), replicated across multiple nodes via a fault-tolerant consensus algorithm.
Four features — note which need consensus and which don't:
| Feature | Needs consensus? | Detail |
|---|---|---|
| Locks and leases | YES | Built on the atomic, fault-tolerant CAS. If several nodes concurrently try to acquire the same lease, only one succeeds |
| Support for fencing | YES | Monotonically increasing ID per log entry — zxid/cversion in ZooKeeper, revision number in etcd |
| Failure detection | no | Clients maintain a LONG-LIVED SESSION and exchange HEARTBEATS. Even if the connection is temporarily interrupted or a server fails, ANY LEASES HELD REMAIN ACTIVE. But if there's no heartbeat for longer than the lease timeout, the service ASSUMES THE CLIENT IS DEAD AND RELEASES THE LEASE (ZooKeeper's ephemeral nodes) |
| Change notifications | no | A client can request notification whenever certain keys change — finding out when another client joins (by the value it writes) or fails (its ephemeral nodes disappear). Saves the client from frequently polling |
Four use cases:
① Configuration management. Store timeouts, thread pool sizes as key-value pairs; processes load on startup and subscribe to changes.
This DOESN'T NEED the consensus aspect, but it's CONVENIENT if you're already running the service. Alternatively, a process could periodically POLL a file or URL, avoiding a specialized service.
② Allocating work to nodes. Choosing a leader/primary among several instances (necessary for single-leader databases, also appropriate for job schedulers and similar stateful systems), and deciding which SHARD to assign to which node — rebalancing as nodes join, taking over as nodes fail.
These can be achieved by judicious use of atomic operations, ephemeral nodes, and notifications. IT'S NOT EASY, despite libraries like Apache Curator — BUT IT IS STILL MUCH BETTER THAN ATTEMPTING TO IMPLEMENT THE CONSENSUS ALGORITHMS FROM SCRATCH, WHICH WOULD BE VERY PRONE TO BUGS.
The key architectural advantage:
A dedicated coordination service can run on a FIXED SET OF NODES (usually three or five), REGARDLESS OF HOW MANY NODES ARE IN THE SYSTEM THAT RELIES ON IT. In a storage system with THOUSANDS of shards, running a consensus algorithm over thousands of nodes would be TERRIBLY INEFFICIENT; it's much better to "OUTSOURCE" THE CONSENSUS to a small number of nodes.
The data-rate constraint:
The data is quite SLOW-CHANGING — "the node running on IP 10.1.1.23 is the leader for shard 7" — changing on a timescale of MINUTES OR HOURS. Coordination services are NOT INTENDED FOR DATA THAT MAY CHANGE THOUSANDS OF TIMES PER SECOND. For that, use a conventional database, or Apache BookKeeper to replicate fast-changing internal state.
③ Service discovery.
Convenient — failure detection and change notification make it easy to track instances as they come and go. And if you're already using it for leases and leader election, it makes sense to use it for discovery too.
HOWEVER, USING CONSENSUS FOR SERVICE DISCOVERY IS OFTEN OVERKILL. This use case GENERALLY DOESN'T REQUIRE LINEARIZABILITY, and it's MORE IMPORTANT THAT IT IS HIGHLY AVAILABLE AND FAST, since without it everything would grind to a halt. It's therefore usually preferable to CACHE service discovery information — clients that can't connect bypass the cache, retry with the latest value, and update it; caches may also refresh on a TTL. (DNS-based discovery uses multiple layers of caching for exactly this reason.)
ZooKeeper OBSERVERS support this: replicas that receive the log and maintain a copy of the data but DO NOT PARTICIPATE IN THE VOTING PROCESS. Reads from an observer are NOT LINEARIZABLE as they might be stale, BUT THEY REMAIN AVAILABLE EVEN IF THE NETWORK IS INTERRUPTED, and they increase read throughput.