etcd is a small, strongly consistent key-value store meant for the data a distributed system cannot afford to get wrong: configuration, service membership, leader identity and, most famously, the entire state of a Kubernetes cluster. It replicates every write through Raft, keeps a versioned history of every key, and lets clients watch that history as a stream. Those three properties make it a coordination service rather than a database, and they also explain most of its operational behaviour.

This article covers how a member is built, revisions, linearizable reads, transactions, watches and leases with Go client code, storage maintenance, a worked recovery, failure modes and a checklist. Version-specific details refer to etcd v3.6. For the consensus algorithm underneath, see Raft consensus, in depth.

Advertisement

What etcd is for, and what it is not

Use etcd when a few processes must agree on a small amount of state and react to changes quickly: leader identity, live nodes, configuration. Writes and default reads are linearizable, and watches deliver every change in order unless history has been compacted.

The price is a Raft majority and an fsync per write, and a data set that must fit on one machine. The documented limits show the intended scale: requests are capped at 1.5 MiB by default (--max-request-bytes), the default storage quota is 2 GiB, and 8 GiB is the suggested maximum for normal environments. etcd is the wrong place for blobs, logs, metrics or anything that grows with user traffic. ZooKeeper occupies the same niche with a different data model; the ZooKeeper deep dive is a useful comparison.

How a member is built

A client speaks gRPC to any member through the KV, Watch, Lease, Cluster, Maintenance and Auth services. Writes are forwarded to the leader, which appends them to the Raft log; each member fsyncs the entry to its write-ahead log before acknowledging, and once a majority has it, every member applies it to its local MVCC store.

That store is bbolt, a single-file B+tree keyed by revision, plus an in-memory treeIndex mapping each key to the revisions at which it changed. A watch hub forwards each applied change to matching watchers.

Inside an etcd member: from gRPC request to watch eventClientclientv3 / etcdctlgRPC APIKV, Watch, Lease, ...Put / TxnRaft (leader)propose entryFollower 1append + fsyncFollower 2append + fsyncWALfsync before acklogApply loopcommitted entriesmajority ackMVCC storerevision++treeIndex (RAM)key to revisionsbbolt backendrevision to KeyValueWatch hubevents by revisionnotifystreamLinearizable read:leader confirms it is stillleader (ReadIndex), memberwaits until applied, thenreads its local MVCC store.Serializable read:any member answers fromlocal state; may be stale.
The write path runs through Raft and the WAL before the change is applied to the MVCC store. Reads, watches and history all hang off the store's revision counter.

After --snapshot-count applied entries (10,000 by default in v3.6) a member snapshots so the WAL can be truncated. Heartbeats default to 100 ms and the election timeout to 1,000 ms; tune both to measured network round-trip time.

Advertisement

Revisions: the key space has a clock

etcd keeps one 64-bit counter for the whole key space: the revision. Every change, whether a put, a delete or a whole transaction, increments it once, and all writes in one transaction share that revision. Each key carries three numbers: create_revision (when the current incarnation of the key was created), mod_revision (when it last changed) and version (how many times it has changed since creation; a delete resets it). Every response carries a header with the revision the cluster was at when it answered.

$ etcdctl put /config/feature-x on
OK
$ etcdctl put /config/feature-x off
OK
$ etcdctl get /config/feature-x --write-out=json | jq '.header.revision, .kvs[0]'
1043
{ "key": "L2NvbmZpZy9mZWF0dXJlLXg=", "create_revision": 1042,
  "mod_revision": 1043, "version": 2, "value": "b2Zm" }      # keys and values are base64
$ etcdctl get /config/feature-x --rev=1042                     # read the past
/config/feature-x
on
$ etcdctl get /config/ --prefix --keys-only                    # "directories" are just prefixes
$ etcdctl get /config/feature-x --consistency=s                # serializable: local, maybe stale

The key space is flat and byte-ordered; prefixes stand in for directories. Old revisions are kept until compaction, so you can read the past with --rev and watch from any retained revision. Kubernetes builds on this directly: an object's resourceVersion comes from etcd's mod revision, and API server updates are compare-and-swap transactions underneath.

Reads: linearizable by default, serializable on request

A linearizable read must see every write that completed before it started, even on a lagging follower or a deposed leader. etcd uses Raft's ReadIndex: the leader confirms its leadership with a majority and returns its commit index, and the member waits until it has applied that index before reading locally. That costs a round trip but no log write.

A serializable read, --consistency=s in etcdctl or clientv3.WithSerializable() in Go, skips that step and answers from whatever the local member has applied. It is faster, keeps working on a member cut off from the majority, and may be stale. Use it for monitoring and caches, never for a decision that a write depends on. The quorum arithmetic is the usual one:

MembersMajorityFailures toleratedNotes
110development only
321the common production size
532survives a failure during maintenance
743rarely worth it; every write waits on more fsyncs

Even member counts add cost without adding tolerance; four members tolerate one failure, like three.

Transactions: compare, then act

An etcd transaction is a single atomic if-then-else: a list of comparisons on keys' value, version, create revision, mod revision or lease, a list of operations to run if all comparisons hold, and another list if any fails. There are no multi-round interactive transactions. This is enough for compare-and-swap, which is the building block for everything else:

import (
    "context"
    clientv3 "go.etcd.io/etcd/client/v3"
)

// Compare-and-swap: update key only if nobody changed it since we read it.
func casUpdate(ctx context.Context, cli *clientv3.Client, key string,
    mutate func([]byte) []byte) error {
    for {
        get, err := cli.Get(ctx, key)
        if err != nil {
            return err
        }
        var old []byte
        rev := int64(0) // ModRevision 0 means "key does not exist"
        if len(get.Kvs) == 1 {
            old, rev = get.Kvs[0].Value, get.Kvs[0].ModRevision
        }
        resp, err := cli.Txn(ctx).
            If(clientv3.Compare(clientv3.ModRevision(key), "=", rev)).
            Then(clientv3.OpPut(key, string(mutate(old)))).
            Commit()
        if err != nil {
            return err
        }
        if resp.Succeeded {
            return nil // committed at revision resp.Header.Revision
        }
        // someone else won the race: loop, re-read, retry
    }
}

If another client wrote the key after the read, its ModRevision differs, the transaction fails and the loop retries; comparing against 0 gives create-if-absent. A transaction may contain at most --max-txn-ops operations, 128 by default, and the whole request is still bound by the request size limit.

Watches: list, then watch from the next revision

A watch is a long-lived gRPC stream of events for a key or prefix, starting at a revision you choose. The safe pattern for building a local cache is to list the prefix, note the header revision of that response, and start watching at the next revision. Nothing can fall in the gap, because the list was a consistent snapshot at exactly that revision:

// List-then-watch, the pattern Kubernetes informers use.
func syncPrefix(ctx context.Context, cli *clientv3.Client, prefix string) error {
    for {
        list, err := cli.Get(ctx, prefix, clientv3.WithPrefix())
        if err != nil {
            return err
        }
        rebuildCache(list.Kvs)
        next := list.Header.Revision + 1 // resume exactly after the snapshot we listed

        wch := cli.Watch(ctx, prefix, clientv3.WithPrefix(),
            clientv3.WithRev(next), clientv3.WithProgressNotify())
        for wresp := range wch {
            if wresp.CompactRevision != 0 {
                break // history we needed is gone: fall through and re-list
            }
            if err := wresp.Err(); err != nil {
                break
            }
            for _, ev := range wresp.Events {
                apply(ev.Type, ev.Kv) // PUT or DELETE, ev.Kv.ModRevision is the event's revision
            }
        }
        if ctx.Err() != nil {
            return ctx.Err()
        }
    }
}

If the start revision was already compacted, the server cancels the watch and reports the compaction revision; the only correct response is to list again. WithProgressNotify makes the server send periodic empty responses with the current revision, so a watcher on a quiet prefix knows how current it is.

Leases, locks and fencing

A lease is a server-side TTL timer kept alive by client keepalives; keys attached to it are deleted when it expires. A service registers /services/api/node-7 under its lease, and if the process dies, watchers see the delete. After a leader change the new leader refreshes outstanding leases rather than expiring them, so a lease can outlive its nominal TTL.

The concurrency package builds mutexes and elections on this. A mutex holder is the key with the lowest create revision under the lock prefix, and each waiter watches only its predecessor:

import "go.etcd.io/etcd/client/v3/concurrency"

sess, err := concurrency.NewSession(cli, concurrency.WithTTL(10)) // lease + keepalive
if err != nil { return err }
defer sess.Close()

m := concurrency.NewMutex(sess, "/locks/reindex")
if err := m.Lock(ctx); err != nil { return err }

// Fence writes inside etcd on still owning the lock:
resp, err := cli.Txn(ctx).If(m.IsOwner()).
    Then(clientv3.OpPut("/jobs/reindex/state", "running")).Commit()
if err == nil && !resp.Succeeded { /* lost the lock: stop */ }

select {
case <-sess.Done(): // lease expired or keepalive failed: we no longer hold the lock
    stopWork()
case <-workFinished:
}
_ = m.Unlock(ctx)

A lock does not stop a paused process from acting after its lease expired. Inside etcd, guard writes with m.IsOwner() in the transaction, as above. For writes to other systems, pass a fencing token, such as the lock key's create revision, and have the resource reject stale tokens. Fencing tokens and leases cover the theory.

Storage: compaction, defragmentation and the quota

Because every revision is kept, the database grows with write traffic, not key count. Compaction discards revisions older than a chosen one, except each key's latest version, manually or via --auto-compaction-mode (periodic or revision) and --auto-compaction-retention; the server default retention is 0, which means auto-compaction is off. Kubernetes clusters normally rely on the API server, which issues compactions on its own schedule.

Compaction frees pages inside the bbolt file but does not shrink the file. Defragmentation rewrites the file to reclaim that space, and while it runs the member does not serve reads or writes, so defragment one member at a time and never the whole cluster at once. Finally, the space quota (--quota-backend-bytes, 2 GiB by default) is a safety stop. When the backend exceeds it, the cluster raises a NOSPACE alarm and accepts only reads and deletes until an operator compacts, defragments and disarms the alarm.

Worked example: recovering from NOSPACE

Scenario: a three-member cluster behind a Kubernetes control plane starts rejecting writes with etcdserver: mvcc: database space exceeded. A controller has been rewriting a large custom resource status every second, and compaction was not keeping up. The recovery is the documented maintenance sequence:

# 1. Confirm the alarm and see sizes and the current revision on every member
etcdctl --endpoints=$EPS alarm list
etcdctl --endpoints=$EPS endpoint status --write-out=table

# 2. Compact away old history (keeps only the latest revision of each key)
rev=$(etcdctl --endpoints=$EP1 endpoint status --write-out=json \
      | grep -o '"revision":[0-9]*' | grep -o '[0-9]*')
etcdctl --endpoints=$EP1 compaction "$rev"

# 3. Defragment ONE member at a time; each is unavailable while it rewrites its file
for ep in $EP1 $EP2 $EP3; do etcdctl --endpoints=$ep defrag; done

# 4. Clear the alarm, then fix the cause (auto-compaction, quota, a runaway writer)
etcdctl --endpoints=$EPS alarm disarm

Then fix the cause: rate-limit the writer, confirm compaction runs, and only then consider a larger quota, staying under 8 GiB. Alert on etcd_mvcc_db_total_size_in_bytes against the quota.

Operating etcd

  • Disk latency is the first metric. Every commit waits for fsync. Watch etcd_disk_wal_fsync_duration_seconds and etcd_disk_backend_commit_duration_seconds at the 99th percentile; dedicated SSDs and no noisy neighbours on the same disk matter more than CPU.
  • Leader stability. etcd_server_has_leader should always be 1, and etcd_server_leader_changes_seen_total should rarely move. Frequent elections usually mean slow disks, CPU starvation or heartbeat settings too tight for the network.
  • Backups. Take regular snapshots with etcdctl snapshot save from one healthy member and test restores. From v3.6, restoring is done with etcdutl snapshot restore; the etcdctl restore command is gone.
  • Membership changes one at a time. Add new members as learners (member add --learner), let them catch up, then member promote. A learner does not count towards quorum, so adding one cannot cost you availability.
  • Security. Enable mutual TLS for peers and clients plus role-based auth; whoever can write a Kubernetes cluster's etcd owns the cluster.

Failure modes

SymptomLikely causeWhat to do
Leader elections every few minutesfsync latency spikes, CPU starvation, timeouts tuned below network RTTMove to dedicated SSD, check disk metrics, raise election timeout in proportion
Writes fail with database space exceededNo or slow compaction, runaway writer, quota too smallCompact, defrag one member at a time, disarm, fix the writer
Watchers resync constantlyCompaction retention shorter than consumer lagRetain more history or make consumers faster; always handle compaction
Writes hang, reads still work on some membersQuorum lostRestore majority; never force a new cluster from one member unless recovering from backup deliberately
Lock held by two processesLease expired during a pauseFence writes with IsOwner or fencing tokens; do not rely on the lock alone
High memory and slow list callsLarge values or huge prefixes read wholeKeep values small, paginate, store bulk data elsewhere

What to do next

  1. Inventory what you store in etcd and move anything large or traffic-proportional elsewhere; keep values well under the 1.5 MiB request limit.
  2. Run three or five members on dedicated SSDs, set heartbeat and election timeouts from measured RTT, and alert on fsync p99 and leader changes.
  3. Confirm something compacts history regularly (auto-compaction or the Kubernetes API server), and alert on database size against the quota.
  4. Schedule rolling defragmentation, one member at a time, and rehearse the NOSPACE recovery on a staging cluster.
  5. Use compare-and-swap on ModRevision for every read-modify-write, and list-then-watch with compaction handling for every cache.
  6. Build locks and leader election on the concurrency package, and fence downstream writes with IsOwner or a fencing token.
  7. Automate snapshot backups and test a full restore with etcdutl at least once a quarter.
  8. Enable TLS and authentication, and restrict who can reach the client port.
Key takeaway: etcd is a Raft-replicated, revisioned key-value store for small, critical coordination state. Its single revision counter powers compare-and-swap transactions, consistent reads of the past and gap-free list-then-watch caches, while leases turn liveness into keys and locks. Keep it healthy by keeping data small, disks fast, history compacted, files defragmented one member at a time, backups tested and every lock fenced.