An Introduction to etcd: The Memory of a Kubernetes Cluster

All of the state in a Kubernetes cluster lives in one place: etcd. For engineers with a backend development background, it is essentially a database. Familiar concepts such as replication, transactions, and optimistic locking all have their counterparts here.

Which core mechanisms make etcd Kubernetes’ single source of truth? And how much do we need to understand as DevOps engineers?

This article is based on etcd v3.5 and Kubernetes v1.37. It skips dense mathematical proofs and source code details to focus on core mechanisms and day-to-day operations.


Why Kubernetes Needs etcd

etcd is a distributed key-value store designed for small amounts of critical data that must remain consistent across machines. Its name combines the Unix configuration directory /etc with d for “distributed.” CoreOS developed it in Go in 2013. Its core mission is to let multiple machines maintain “one truth everyone agrees on” and to notify clients as soon as that data changes.

etcd’s purpose comes into focus when set against familiar backend development pain points:

etcd’s answer is distributed consensus: a write succeeds only after a majority of cluster members—a quorum—have successfully persisted it to disk. This rule provides both availability, by tolerating failures in a minority of nodes, and consistency, by ensuring there is always only one valid state.

Around this consensus core, etcd builds four supporting mechanisms:

  1. The Raft consensus algorithm determines which writes and proposals are accepted.
  2. MVCC and Revision use multi-version concurrency control to retain historical versions of every data change.
  3. Watch lets clients subscribe to changes and replay events from a specified historical revision.
  4. Lease provides a foundation for time-to-live (TTL) expiration and renewal through heartbeats.

Kubernetes uses etcd because every control plane decision must be based on current, authoritative state. Pods, Nodes, Deployments, and Secrets are all encoded with Protobuf and stored under etcd’s /registry/ prefix.

Within a Kubernetes cluster, only kube-apiserver reads and writes etcd directly, using the --etcd-servers setting. For high event volumes, --etcd-servers-overrides can also route Events to a separate etcd cluster.

In terms of the CAP theorem, etcd favors CP: consistency and partition tolerance. During a network partition, it would rather stop accepting writes on the minority side than return potentially stale, inconsistent data.

This design makes etcd a store for cluster metadata, rather than a general-purpose database. Its default backend quota is just 2 GiB (--quota-backend-bytes), and the official suggestion for normal environments is a maximum of 8 GiB. Every write must complete the consensus process. The upper bound on write latency is determined by how long the slowest node in the majority takes to write the data to disk.


Raft Consensus: How Majority Agreement Provides Reliability

Raft is the distributed consensus algorithm used by etcd. It was designed to make consensus more intuitive and easier to understand than the established Paxos algorithm. To understand how Raft behaves in day-to-day operations, you need one core principle, three roles, and two key processes.

 [Raft Write Flow]

 Client ---> [ Leader ]
             |
             |--(1) Append Log Entry--> [ Follower 1 ]
             |                          (Disk OK)
             |
             |--(2) Append Log Entry--> [ Follower 2 ]
             |                          (Disk OK)
             |
             v
       (Quorum Reached: 2/3)
             |
             |--(3) Commit & Apply Machine
             v
 Client <--- [ Response OK ]

One Core Principle

Every state change requires agreement from a majority of members—a quorum. The key is that only one side can have a majority. If a three-node cluster is partitioned into groups of two and one, only the two-node side can form a majority and continue accepting writes. The isolated node cannot complete any write. The two sides cannot both hold a majority at the same time, so they cannot produce two conflicting versions of the data. The rule itself prevents split-brain.

Three Roles and Two Processes

At any given time, the cluster has at most one Leader, with the other nodes acting as Followers. When a Follower starts an election, it temporarily becomes a Candidate. The Leader handles all write requests.

For reads, etcd defaults to linearizable reads. Even when a request goes to a Follower, that node first checks the latest committed progress with the Leader, ensuring the read returns the latest committed data. For maximum throughput, clients can instead use local reads, known as serializable reads. A single node answers directly from its local state, which is faster but may return briefly outdated data.

Fault-Tolerance Arithmetic and Why Node Counts Should Be Odd

For a cluster of NN nodes, the number of failures it can tolerate is:

⌊(N−1)/2⌋{\lfloor(N-1)/2\rfloor}

In other words, divide N−1 by 2 and round down. A three-node cluster can tolerate one failure; a five-node cluster can tolerate two.

Total nodes (NN)Nodes required for quorumTolerated node failuresFault-tolerance explanation
22 (agreement from all 2/2 nodes)0One failure prevents a majority; no fault tolerance
32 (a 2/3 majority)1Two nodes can continue serving after one failure
43 (a 3/4 majority)1Same fault tolerance as three nodes, with unnecessary extra cost
53 (a 3/5 majority)2Recommended configuration for production

The mathematics of majority agreement shows that two nodes tolerate zero failures, while four nodes provide exactly the same fault tolerance as three, with additional communication overhead and opportunities for failure. An etcd cluster should therefore have an odd number of members, usually three or five.

DevOps engineers need a clear understanding of how nodes on the minority side behave. When the cluster loses quorum, etcd immediately rejects all writes. Linearizable reads, the default, also fail. Only local reads can still respond, and they may return stale data.

For Kubernetes, if two nodes in a three-node etcd cluster fail, the entire control plane freezes: resources cannot be created or modified, and the scheduler cannot assign new Pods. However, application Pods already running on worker nodes remain completely unaffected and continue handling network traffic normally.


MVCC and Revision: Every Change Becomes History

In traditional key-value databases, a PUT operation usually overwrites data in place, erasing the old value. etcd uses multi-version concurrency control (MVCC): each write adds a new version instead of overwriting the old one, and assigns it a globally increasing Revision.

Revision is a 64-bit, cluster-wide counter. Regardless of how many keys a transaction changes, the global Revision increases by just one. It acts as a global logical clock for the entire distributed system.

 [MVCC Key Lifecycle]

 Global Revision:  1          2          3          4
                   |          |          |          |
 Key: /app/config  |-- "v1" --|-- "v2" --|-- "v3" --|
                   |          |          |
                   (create=1) (mod=2)    (mod=3)
                   (ver=1)    (ver=2)    (ver=3)

Notice that the global Revision has reached 4, while /app/config only extends to 3. Revision 4 comes from a write to another key. This is what makes the counter “global”: every write advances it, regardless of which key changes.

Querying a key in etcd returns three core attributes:

$ etcdctl put /app/config "v1"
$ etcdctl put /app/config "v2"
$ etcdctl get /app/config -w json

Optimistic Concurrency Control and Event Replay

Retaining version history gives Kubernetes two key capabilities:

  1. Optimistic concurrency control (CAS, Compare-And-Swap) is what backend developers commonly call optimistic locking. Instead of acquiring a lock before writing, the client attaches a condition to the update: “Allow this write only if the data is still at the version I read.” If someone else has changed it in the meantime, the write fails, and the client reads again before retrying. A common backend implementation adds a version column to a table and checks it with UPDATE ... WHERE version = ?. etcd’s conditional transactions (If / Then / Else) do the same thing, with mod_revision serving as that version field. A client can require that an update proceed “only if mod_revision == X.” The resourceVersion field in Kubernetes objects is based on this mechanism and maps directly to the object’s mod_revision in etcd. When kubectl edit or a controller updating an object encounters a 409 Conflict (“the object has been modified”), the underlying cause is a failed optimistic locking check.
  2. Event replay, the foundation of Watch. Because the history is retained in full, a client can request at any time: “Replay all change events after Revision 100, in order.”

History Compaction

If history grows indefinitely, the disk will eventually fill up. etcd uses compaction to remove outdated history:

$ etcdctl compact 5
$ etcdctl get /app/config --rev=4
# Output: rpc error: ... etcdserver: mvcc: required revision has been compacted

After compaction, any request attempting to read or watch a compacted Revision receives an ErrCompacted error.

For automated operations, etcd supports --auto-compaction-mode, which can be configured for periodic compaction or revision-based compaction. In Kubernetes, kube-apiserver sends compaction requests to etcd every five minutes by default (--etcd-compaction-interval=5m0s).

One detail matters here: compaction does not shrink the database file on disk. The defragmentation section below explains why and how to reclaim that space.


Watch and Lease: Events and Heartbeats

Watch: A Reliable Event Stream

A client can specify start_revision when establishing a Watch. etcd then pushes all events that occurred after that revision to the client in order.

$ etcdctl watch /app/config --rev=6

This design handles disconnections well. When a client reconnects, it only needs to supply the last Revision it received. etcd first replays the events missed during the disconnection, then seamlessly resumes live notifications.

Kubernetes’ controller client library, client-go, provides a caching and event synchronization mechanism called an Informer. Built-in controllers, custom Operators, and kubectl get -w all synchronize state with the apiserver through the List-and-Watch pattern:

 [Informer List-and-Watch Pattern]

 Controller / Informer
    |
    |--(1) List Request (Full State)
    |      ---> kube-apiserver ---> etcd
    |
    |<-(2) Return Objects + resourceVersion=100
    |
    |--(3) Watch Request (resourceVersion=100)
    |      ---> kube-apiserver ---> etcd
    |
    |<-(4) Stream Incremental Events (rev > 100)

Clients connect to kube-apiserver rather than directly to etcd. kube-apiserver has a built-in Watch Cache (--watch-cache, enabled by default) that handles large numbers of Watch connections and caches events locally.

If a controller remains disconnected long enough that etcd has already compacted away its requested Revision, kube-apiserver returns 410 Gone (“too old resource version”) to the Watch. The controller then automatically performs a re-list, fetching the full data again and establishing a new Watch.

Lease: Managing Key Lifetimes with TTL

A Lease is an etcd primitive for managing key lifecycles. A client first requests a Lease ID with a time to live (TTL), then attaches specific keys to that lease.

The client must periodically send KeepAlive heartbeats to renew the lease. If the client crashes and stops renewing, once the lease expires, etcd automatically deletes every key attached to that lease and emits deletion events. If you are developing microservices and using etcd clientv3 for service registration and health checks, Lease is the essential underlying tool.

Clearing Up a Misconception: Kubernetes Node Heartbeats

Many engineers assume that Kubernetes Node heartbeats use native etcd leases. They do not.

Because only kube-apiserver can access etcd directly, Kubernetes implements node heartbeats at the API object layer. A separate Lease object is much lighter than updating the entire NodeStatus each time, substantially reducing heartbeat overhead in large clusters:


Four Practical Operational Scenarios

Understanding the underlying principles makes it easier to identify the cause of an unexpected production failure. The following are four of the most common scenarios DevOps engineers encounter when operating etcd.

1. Backend Quota Exhaustion: The Cluster Becomes Read-Only

etcd has a default quota limit of 2 GiB. Once the database file exceeds that quota, etcd immediately raises a NOSPACE alarm and rejects all writes, allowing only reads and deletes. In Kubernetes, this leaves the entire cluster unable to create Pods or update any state.

The most common cause is an event storm: application Pods stuck in crash-and-restart loops, or repeated scheduling failures, generate enough Events to fill the storage space.

# Four-step recovery procedure:
# 1. Get the current Revision
$ rev=$(etcdctl endpoint status --write-out="json" | jq .[0].Status.header.revision)

# 2. Compact outdated history
$ etcdctl compact $rev

# 3. Reclaim fragmented space on the filesystem
$ etcdctl defrag

# 4. Clear the NOSPACE alarm to restore writes
$ etcdctl alarm disarm

Operators should routinely monitor the following Prometheus metrics:

The difference between them is internally fragmented space.

2. Defragmentation: Why Compaction Does Not Shrink the File on Disk

etcd uses bbolt, a B+Tree-based key-value engine, as its storage backend. Compaction and deletion only mark database pages as free internally, making them available for reuse by future writes. The file size visible to the operating system does not shrink.

Running etcdctl defrag rebuilds the database file and returns space to the operating system. However, defragmentation temporarily blocks reads and writes on the node where it runs. In production, never defragment the entire cluster simultaneously. Run it one node at a time, preferably outside peak traffic hours.

3. Cascading Latency: How a Disk IOPS Bottleneck Spreads Failures

Raft requires every write to receive confirmation that a majority of nodes have flushed it to disk with fsync. If a node writes to disk slowly, it slows overall API responses, and the latency propagates upward to kube-apiserver and all downstream controllers.

A more serious outcome occurs when heavy I/O contention causes the etcd Leader to keep missing its 100ms heartbeats, long enough that Followers’ election timeouts expire. Followers then conclude the Leader is lost and start a new election. Repeated elections can leave the cluster rejecting writes for an extended period.

Deployment principles:

4. Backup and Restore: The Last Line of Defense

If a majority of etcd nodes are permanently damaged and cannot be recovered, having no backup means losing the cluster state entirely.

In the etcd toolchain, online operations such as saving snapshots use etcdctl. Offline data directory maintenance and restoration have been handled by the separate etcdutl tool since etcd v3.5:

# 1. Create an online snapshot
$ etcdctl snapshot save /backup/etcd-snapshot-$(date +%Y%m%d).db

# 2. Check snapshot status (inspect offline with etcdutl)
$ etcdutl snapshot status /backup/etcd-snapshot-*.db

The standard restore procedure is:

  1. Stop all kube-apiserver instances to prevent inconsistent data from being written during restoration.
  2. Use etcdutl snapshot restore on every etcd node to rebuild a fresh data directory.
  3. Start all etcd nodes and verify cluster health.
  4. Restart kube-apiserver, kube-scheduler, and kube-controller-manager.

Mechanisms at a Glance

The following table summarizes the article’s core mechanisms and how Kubernetes uses them:

Core etcd mechanismEssential ideaKubernetes mapping and application
Raft consensusWrites require persistence by a majority; reads are linearizable by defaultUse an odd number of members, three or five; losing quorum freezes the control plane while existing Pods continue running
MVCC + RevisionA globally increasing logical clock for version historyObject versions through resourceVersion; optimistic locking and conflict detection through 409 Conflict
WatchSubscribe from a specified Revision and replay a stream of incremental eventsInformer List-and-Watch and the kube-apiserver Watch Cache
LeaseA declaration of liveness with TTL and KeepAliveNode status reporting through an API object layer and lock acquisition for control plane leader election

Understanding etcd’s core principles gives us a clear mental model for tracing failures when building and maintaining Kubernetes clusters: choosing a node count that provides the required fault tolerance, ensuring adequate disk I/O performance, and responding methodically when the backend quota fills up.