← All writing
articleJan 12, 202419 min read

Distributed Configuration Management: The Part Where Everyone Builds a ZooKeeper

Raft, linearizable reads, watch fan-out, and MVCC — how etcd and Consul actually keep fifty thousand services watching config without falling over.

etcdRaftConsensusArchitecture
Distributed Configuration Management: The Part Where Everyone Builds a ZooKeeper cover illustration

Every organisation reaches the same point. Config was in Git. Git was in a deploy pipeline. Then a service needed to change a flag without a redeploy, so someone added a database table, then an admin UI, then a cache, then a notification mechanism, and then they built a distributed systems problem and called it a config service.

The reason it goes wrong is not that config is hard. It is that the requirements look tiny and the failure modes are not.

Store some keys, watch for changes, never serve a stale value, expire things automatically. Five sentences. Then the real system has fifty thousand watch connections, a leader that fails every few weeks, and a compaction policy that silently breaks every client that was offline too long.

The scale that decides the design

A platform serving ten thousand service instances, each watching about five configuration keys.

10,000 instances x 5 watched keys      = 50,000 active watches
one config change must reach them all  within 100ms
that is 50,000 gRPC stream writes      in a tenth of a second
                                       = 500,000 writes/sec fan-out

stored data
100,000 keys x ~1 KB                   = 100 MB
comfortably in RAM on every node

write rate
100 config changes/sec                = high for a config system
(every deploy, every flag flip, every
 auto-scaling policy evaluation)

read rate
5,000 reads/sec
(services fetching config at startup, and on watch miss)

Two observations from those numbers that should drive everything.

The data is tiny and the notification load is enormous. A hundred megabytes replicated three ways is nothing. Fifty thousand concurrent subscriptions and the fan-out that goes with them is the actual system. A design that optimises the storage and treats watches as a bolt-on has optimised the easy part.

A hundred writes per second is not a throughput problem, it is a latency problem. Raft serialises writes through a single leader, and there is no way around that without giving up consistency. So the question is not “can we do a hundred writes a second” — trivially yes. The question is “what is the p99 of a single write, and what does it cost in fsync.”

write path, per committed entry
  leader fsync to WAL              ~0.5ms on NVMe
  round trip to 2 followers         ~0.5ms same rack
  followers fsync (in parallel)      ~0.5ms
                                    ------
  best case                         ~1.5ms
  real p99 under load                5-20ms

On network-attached storage the fsync alone is 5-20ms and everything gets worse by an order of magnitude. This is why the first operational rule of etcd is “put the WAL on local NVMe, never on NFS.”

Everything goes through the leader

The architecture is a replicated state machine. Three or five nodes, one leader, Raft consensus.

        client writes
             |
             v
     +---------------+
     |    LEADER    |  accepts all writes
     |  term 8      |  appends to its log
     +---+-------+--+
         |       |
    AppendEntries AppendEntries
         |       |
     +---v--+ +--v----+
     | F1   | |  F2   |   followers
     +-----+ +-------+
         |
    quorum = 2 of 3
    entry is committed once
    a majority has it on disk
         |
         v
   apply to the state machine
   (the in-memory B-tree)

Raft’s job is to make every node apply the same sequence of commands in the same order. Everything else is built on that guarantee.

Time is divided into terms. Each term begins with an election. A follower that has not heard from the leader within a randomised election timeout, typically 150-300ms, assumes the leader is dead, increments its term, and starts asking for votes.

A node grants at most one vote per term, and only to a candidate whose log is at least as up to date as its own. That second condition is the part that makes Raft safe: a candidate missing committed entries can never win, so committed entries can never be lost.

Once a candidate has a majority of votes it becomes leader and starts sending heartbeats as empty AppendEntries calls. The heartbeats do two jobs: they suppress new elections, and they let followers reconcile any entries they missed.

The client write path is then:

  1. Leader receives the write, assigns it an index and the current term.
  2. Leader appends to its own log and fsyncs before acknowledging.
  3. Leader sends AppendEntries to all followers.
  4. Followers append, fsync, and acknowledge.
  5. When a majority has acknowledged, the entry is committed.
  6. Leader applies it to the state machine and replies to the client.

Step 2 is where durability comes from and step 5 is where consistency comes from. A leader can acknowledge a write that later turns out not to be committed, which is why the client’s retry needs to be idempotent — the same write arriving twice must not be applied twice.

The critical section is the fsync, and it is why cluster size is capped. A five-node cluster needs three acknowledgements, a three-node cluster needs two. Every follower you add past three costs write latency and buys almost nothing, because your realistic failure tolerance is dominated by correlated failures (a rack, an AZ, a bad firmware rollout) rather than independent ones.

The read problem is the interesting problem

A read is easy if you don’t care about freshness. Read from the local B-tree, done.

A linearizable read has a much stronger promise: the read observes the most recent committed write, as if it happened immediately before the read. That is what a system providing configuration needs, because “read your own write” across a client that just did a write is the minimum a caller can expect.

And it is not free. If a node is a follower, and the leader commits a write, and then this follower answers a read, it returns a stale value — while looking perfectly healthy.

There are three ways to fix this and the choice is a real trade-off.

Serve reads from the leader only. Correct, trivially so, and it puts every read on the leader’s critical path. Simple, and the leader becomes a throughput ceiling.

Read from followers with an acknowledged commit index. A follower tells the leader “I have everything up to index N.” The leader confirms N is still committed and replies. The follower then serves reads at or below N. Fast, and the follower can be arbitrarily stale — the classic “your config is four seconds old” bug.

ReadIndex, which is what etcd does. Before serving the read, the leader confirms it is still the leader by getting an acknowledgement from a majority, then waits for its state machine to catch up to the confirmed commit index, then reads locally. No log round trip, no waiting for a new entry to be written, and the read is linearizable.

ReadIndex
1. leader records current commit index
2. leader sends a heartbeat round
3. a majority acknowledges
   -> if the leader had been deposed,
      the new leader's term would be higher
      and this fails
4. leader waits for apply queue to reach that index
5. leader serves the read from local state

The cost is one network round trip per linearizable read. That is the whole price, and it is why etcd defaults to it.

The temptation to “optimise” this away is the single most common way to build a config system that eventually serves wrong answers. A lease-based read optimisation — leader reads are served without a round trip if it can prove the lease has not expired — is sound in theory and dangerous in practice, because the proof depends on clock behaviour across nodes. If the clock on a partitioned node looks fine, it will serve stale reads for the remainder of the lease.

Watch is where this actually gets hard

A watch is a long-lived stream from a client to a server. The client says “send me everything that happens under /prod/payments/” and the server keeps the stream open indefinitely.

Server side, that means a registry mapping prefixes to live watchers.

+-----------------------------------------------+
| watch registry                                |
|                                               |
|  /prod/payments/     ->  [w1, w2, w7, ...]   |
|  /prod/checkout/     ->  [w3, w9]            |
|  /prod/feature-flags ->  [w1, w4, w11, ...]  |
+-----------------------------------------------+

On a commit, the server walks the registry, finds every watcher whose prefix the key falls under, and writes an event to each watcher’s stream.

This is O(watchers per prefix) per write, and it is the scaling limit of the entire system.

Worse, and this is the part people miss: a change to a single key under a watched prefix fans out to all watchers of that prefix, not just the ones watching that key. A watch on /prod/payments/ is a subscription to a range, so writing one key under it notifies every client watching that range even though only one key changed.

Now the arithmetic. One config change, 2,000 watchers on the affected prefix, and a 100ms delivery budget.

2,000 stream writes
in 100 ms
= 20,000 writes/sec from a single thread

Each of those writes is a syscall, a buffer copy, and a TLS record. Done naively this is a garbage collection problem before it is a network problem. The mitigations, in order of how much they help:

Do not do the fan-out on the commit path. The commit thread’s only job is to durably record the entry and apply it to the state machine. Push the event onto a channel and let a dedicated watch server consume it. Commit latency and watch latency become independent, which is the entire point.

Batch events per watcher. If a watcher has forty queued events, write one frame containing forty events. Under a config push that touches hundreds of keys, this collapses thousands of stream writes into a handful.

Batch the syscall. Group the writes across watchers so they land in fewer writev calls rather than one call per watcher.

Coalesce. If a watcher has events for the same key, only the last one matters. A client that is three config versions behind does not need to apply all three in order; it needs the current value.

Decouple the slow watcher. A client on a bad connection can block the fan-out for everyone. Every real implementation has a per-watcher outbound buffer with a bounded size, and when the buffer fills, the honest thing to do is terminate that watch and let the client reconnect and resync from a revision. Dropping a slow client’s events without telling it is how you get silent config divergence.

That last point is the important design decision. The two options are:

slow client, buffer full

  option A: drop events, keep the stream
            -> client is now permanently out of date
            -> divergence with no signal

  option B: close the watch
            -> client reconnects
            -> asks for events from its last known revision
            -> resyncs from history

Option B is correct. It costs a reconnect and it turns a silent correctness bug into a visible, self-healing event.

MVCC: one decision that buys three features

Every write creates a new revision rather than overwriting. A key has a history.

revision   key                        value
   1       /prod/db/host              db-1.internal
   2       /prod/db/host              db-2.internal
   3       /prod/payments/timeout     30s
   4       /prod/db/host              db-3.internal

Keys are stored in a B-tree indexed by (key, revision). The “current” value is simply the highest revision for that key. Old revisions are removed by compaction, which keeps the most recent N revisions or everything newer than a time threshold.

That single mechanism gives you three things that would otherwise each need their own design.

Watch catch-up. A client that was disconnected knows the last revision it saw. It asks for events from that revision and gets exactly what it missed. No polling, no full refetch, and the server does not need to know the client exists between connections.

Optimistic concurrency. A transaction can say “modify this key only if it is still at revision 4.” That is compare-and-swap, expressed as a precondition on history rather than a lock.

A fencing token, free. The revision is a globally monotonic number. Every successful write gets one, which makes it a natural fencing token for anything the write protects — the leader election case, the lock service case, the “has my session already been superseded” case.

Without MVCC, catch-up requires the client to refetch everything, CAS requires a version column, and fencing requires a separate counter that can drift out of sync with the store. The version design is what makes etcd a coordination primitive rather than just a config store.

Compaction is where this bites. Compact aggressively and a client that was offline for longer than the retention window cannot catch up, so it must do a full refetch. That is correct behaviour but it must be a designed path, not an error. Conservative compaction preserves more history at the cost of memory that grows without bound, and the system that ran fine for a year starts getting OOM-killed because nobody revisited a setting that was fine when there were a thousand keys.

Leases: a key-value store becomes a coordination system

A lease is a TTL attached to keys. While the lease is alive the keys exist. When it expires, the keys are deleted.

client                          leader
  |  LeaseGrant(ttl=10s)         |
  |----------------------------->|   -> lease_id = 7734
  |<-----------------------------|

  |  Put(/locks/leader, me)      |
  |    with lease_id=7734        |
  |----------------------------->|   committed at revision 912
  |<-----------------------------|

  |  LeaseKeepAlive(7734)        |
  |----------------------------->|   every 3 seconds
  |                              |
  |  ...client crashes...        |
  |                              |
  |            10 seconds pass   |
  |                              |   -> lease 7734 expires
  |                              |   -> /locks/leader deleted
  |                              |   -> revision 940: DELETE

The renewal interval has to be comfortably shorter than the TTL. A 10-second lease renewed every 3 seconds tolerates two consecutive renewal failures before the key is at risk. Renewing every 9.5 seconds against a 10-second TTL means one dropped packet loses the lease.

The primitive this builds is the whole reason distributed config stores are used for coordination at all.

Leader election. Every instance tries to create /locks/leader with a 10-second lease. Exactly one succeeds, because key creation is atomic. That instance is the leader. It renews. If it crashes, the lease expires in at most 10 seconds and the key is deleted, and every instance watching that key is woken by the delete event and one of them wins the next attempt.

Service registration. A service writes its address into /services/payments/node-3 with a 15-second lease, renewing every 5 seconds. Consumers watch the /services/payments/ prefix and get the current membership as a live stream. When a node dies and stops renewing, its entry disappears and consumers see the removal as a watch event. That is a complete service discovery system on top of a key-value store, which is a large part of why Consul exists.

The failure mode to understand is the false expiry: a slow network or a stop-the-world GC pause long enough to miss renewals, and a healthy leader voluntarily gives up leadership. This is why lease TTLs are chosen for the worst-case pause, not the median, and why leader election in a healthy cluster should be rare. If you are seeing frequent elections, either the TTL is too aggressive or something on the leader is stalling.

A note on transactions

A transaction is a set of compare operations and two operation sets. All the compares are evaluated against the current state, and either the success set or the failure set runs, atomically.

Txn(
  compare: [
    /locks/leader does not exist,
    /config/version = 42
  ],
  success: [
    put /locks/leader = me, with lease,
    put /config/version = 43
  ],
  failure: []
)

Both evaluate, one of the two sets runs. The evaluate-then-apply is the whole mechanism, and it is why you can implement “create the key if nobody has it” atomically without a separate lock primitive.

The honest limit: transactions touch a bounded number of keys. This is not a general-purpose ACID store and it is not trying to be one. A transaction that reads and writes thousands of keys will hold locks and stall the leader, and the right answer is always to restructure the schema so the critical section is small.

Failure stories worth testing

Kill the leader mid-write

Confirm the entry either committed on a majority or did not commit at all. Never half-committed. Then confirm a new leader is elected within the election timeout and clients are redirected.

Kill two of three nodes

The surviving node cannot commit and cannot serve linearizable reads, and it should say so. Most important: it must refuse rather than silently serve stale reads. This is the test that catches people who removed the ReadIndex.

Partition the leader from a majority

The leader keeps accepting writes and none of them commit. This is the split-brain case, and the answer is that the isolated leader stops being able to serve linearizable reads after its lease expires. A system that keeps serving reads from the isolated side has given up linearizability without meaning to.

Disconnect every watcher for ten minutes, then reconnect

Each client should resync from its last revision and end up correct. If compaction ran during the gap, the client must be told clearly and fall back to a full refetch. Both paths must work.

Push 500 config keys in one commit

This is the fan-out stress test. Watch delivery must stay inside the budget through batching, and the commit latency must not move.

Open a watch on a very broad prefix

Watch / on a busy cluster. It is legal and it will be slow, and the platform should be able to handle it without falling over. If a single client can DoS the watch server, the per-watcher buffer bounds and the disconnect policy are wrong.

Stop renewing a lease

The key should disappear at TTL and watchers should be notified. Then the reverse: keep renewing from a process that has been paused for 20 seconds and confirm the system does not end up with two leaders.

Run compaction with clients mid-resync

Confirm the “revision has been compacted” path returns an explicit error rather than silently returning nothing.

Force a leader election while fifty thousand watches are open

Every watch survives the failover and reconnects to the new leader. This is the reconnect storm, and it is the single most stressful thing that happens to this system in production.

A production-ready architecture

   10,000 clients
        |  gRPC, 50,000 long-lived watch streams
        v
  +---------------------------+
  |  load balancer            |  consistent hash on client id,
  |  (L4)                     |  so a client lands on one node
  +------------+--------------+
               |
     +---------+---------+------------------+
     |                   |                  |
     v                   v                  v
  +--------+        +--------+         +--------+
  | member |        | member |         | member |   Raft: one leader,
  |        |<------>|        |<------>|        |   two followers
  |  raft  | AppendEntries    |  raft  |   each ~100 MB
  +---+----+        +----+----+         +---+----+
      |                  |                 |
      |  apply to B-tree, index (key, revision)
      |                  |
      v                  v
  +----------+     +----------+
  |  WAL     |     | snapshot |
  |  NVMe    |     | periodic |
  +----------+     +----------+

  leader commit thread
        |
        | in-memory channel (non-blocking)
        v
  +---------------------------+
  | watch server              |  registry: prefix -> watchers
  |  batch, coalesce          |  per-watcher bounded buffer
  |  disconnect slow clients  |  drop buffer full -> force resync
  +---------------------------+

  lease manager
  background expiry, key deletion on TTL

A sensible delivery checklist:

  1. Run three or five members. Do not run seven; you pay quorum cost for no real gain.
  2. Put the WAL on local NVMe. Never on NFS, never on a network volume. Measure fsync before anything else.
  3. Serve linearizable reads by default, and measure what the ReadIndex round trip costs you before optimising it away.
  4. Move watch fan-out off the commit path entirely. Commit latency and notification latency are different problems.
  5. Batch and coalesce watch events per watcher, and give every watcher a bounded outbound buffer.
  6. On a full buffer, close the watch rather than dropping events silently.
  7. Use MVCC revisions for watch catch-up, for CAS preconditions, and for fencing tokens.
  8. Design the compaction path deliberately, including the compacted-revision error and the full-refetch fallback.
  9. Set lease TTLs from the worst-case process pause, and renew at well under a third of the TTL.
  10. Build the reconnect path and load-test it, because a leader change with fifty thousand watches is the real peak.
  11. Make every client operation idempotent, because a retried write on a leader that did not commit is unavoidable.
  12. Instrument fsync latency, leader election count, watch delivery latency, watch reconnect rate, and the count of watches dropped for being slow.

Common mistakes

Mistake What actually happens Better decision
WAL on network storage fsync goes from 0.5ms to 10ms+, writes become seconds Local NVMe, measured
Seven or more members Quorum cost rises, latency rises, failure tolerance barely improves Three or five
Watch fan-out on the commit path Commit latency becomes notification latency and both degrade Dedicated watch server, async channel
One stream write per event Syscall and TLS overhead dominate under fan-out Batch events per watcher into one frame
Dropping events for a slow client Client is silently stale forever Close the watch, force resync from revision
Disabling linearizable reads “for speed” Callers silently observe stale config after a leader change Keep ReadIndex, measure the cost
Aggressive compaction Offline clients get an unrecoverable error path Design the compacted-revision fallback
Lease TTL too short Healthy leaders lose leadership during GC pauses Size for worst-case pause, renew under TTL/3
Lease renewal interval near the TTL A single dropped packet loses the lease Renew at a third of the TTL or less
No idempotency on writes A retry after an ambiguous commit applies twice Idempotent keys, CAS on revision
Large transactions Leader stalls, availability drops Restructure so the critical section is small
One watch prefix per key Registry and fan-out scale with key count, not need Watch the range you actually use
No client-side backoff Leader change becomes a reconnect storm Jittered exponential backoff on reconnect
Treating config store as general database Transactions and large key sets stall the leader It is a config and coordination store

The complete story in one minute

A distributed config store is a replicated state machine built on Raft. Three or five nodes elect a leader, all writes go through it, it fsyncs to a write-ahead log before acknowledging, and an entry is committed once a majority has it on disk. The critical path is the fsync plus one network round trip, which is why the WAL lives on local NVMe and the cluster stays small.

Reads are where the subtlety is. A follower can serve a stale answer while looking perfectly healthy, so a linearizable read needs a leadership check first. ReadIndex does that with one heartbeat round trip and no log write, and it is the reason to never quietly “optimise” linearizable reads away.

The watch mechanism is the real system. Fifty thousand long-lived streams, and a single commit fans out to every watcher of the affected prefix, so the fan-out runs off the commit path, batches events per watcher, and gives each watcher a bounded buffer. When a buffer fills, the watch is closed rather than silently starved, because a client that misses events with no signal is worse than a client that reconnects and resyncs from its last revision.

MVCC is what makes that resync cheap, and it also gives compare-and-swap preconditions and a monotonic fencing token for free. Compaction is the trade: aggressive compaction frees memory but bounds how long a client can be offline, so the error path for a compacted revision has to be designed rather than discovered.

Leases are what turn a key-value store into a coordination system. A key attached to a TTL disappears when the client stops renewing, which is leader election and service registration with no extra machinery — at the cost of false expiries when the TTL is chosen for the median pause instead of the worst one.

So the path is:

write -> leader -> fsync -> quorum -> commit
                                       |
                     apply to B-tree ---+--> watch server -> 50,000 streams
                                       |
                     revision recorded -+--> CAS preconditions
                                       |
                     lease attached ----+--> expiry, election, discovery

The hard part was never storing a key. It was making fifty thousand clients that are allowed to be offline believe they are seeing the current value.

Technical references

Keep reading
Browse everything