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:
- A single-node database provides full ACID transaction guarantees. But if its server suffers a hardware failure or crashes, the entire service goes down, creating a single point of failure (SPOF).
- Asynchronous primary-replica replication adds replicas for high availability, but writes that have not yet been replicated can be lost during primary failover. If a network partition occurs, nodes on both sides may independently accept writes, causing split-brain with no way to arbitrate between them.
- Caching systems have similar limitations. Redis’ default replication mode is also asynchronous, so it cannot safely serve as the single source of truth for a distributed system.
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:
- The Raft consensus algorithm determines which writes and proposals are accepted.
- MVCC and Revision use multi-version concurrency control to retain historical versions of every data change.
- Watch lets clients subscribe to changes and replay events from a specified historical revision.
- 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.
- Process 1: Leader election. The Leader sends heartbeats to all Followers at a fixed interval, 100ms by default in etcd. If a Follower receives no heartbeat within the election timeout—1000ms by default, with randomized variation to avoid tied votes—it increments its term, becomes a Candidate, and calls for votes. A candidate that receives a majority of votes becomes the new Leader. During the brief election period, the cluster cannot accept any writes, and client requests time out. This is often the source of the timeout errors commonly seen when the cluster is unstable.
- Process 2: Log replication. When the Leader receives a write request, it first records the operation in its own log and broadcasts it to all Followers in parallel. The write is committed only after a majority of nodes, including the Leader, have successfully flushed the log to disk with
fsync. It is then applied to the state machine—the actual key-value data served to clients—and the client receives a response.
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 nodes, the number of failures it can tolerate is:
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 () | Nodes required for quorum | Tolerated node failures | Fault-tolerance explanation |
|---|---|---|---|
| 2 | 2 (agreement from all 2/2 nodes) | 0 | One failure prevents a majority; no fault tolerance |
| 3 | 2 (a 2/3 majority) | 1 | Two nodes can continue serving after one failure |
| 4 | 3 (a 3/4 majority) | 1 | Same fault tolerance as three nodes, with unnecessary extra cost |
| 5 | 3 (a 3/5 majority) | 2 | Recommended 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:
create_revision: the global Revision at which the key was originally created.mod_revision: the global Revision at which the key was last modified.version: the cumulative number of modifications to the key during its current lifecycle.
$ 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:
- 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
versioncolumn to a table and checks it withUPDATE ... WHERE version = ?. etcd’s conditional transactions (If / Then / Else) do the same thing, withmod_revisionserving as that version field. A client can require that an update proceed “only ifmod_revision == X.” TheresourceVersionfield in Kubernetes objects is based on this mechanism and maps directly to the object’smod_revisionin etcd. Whenkubectl editor a controller updating an object encounters a 409 Conflict (“the object has been modified”), the underlying cause is a failed optimistic locking check. - 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:
- Node heartbeats are maintained in Lease objects (
coordination.k8s.io) in thekube-node-leasenamespace. - kubelet periodically renews its corresponding Node Lease object, every 10 seconds by default.
- The control plane’s Node Lifecycle Controller checks nodes periodically. If a node’s heartbeat is overdue, it marks the node as
NotReadyand triggers Pod eviction after the grace period has elapsed. - The same mechanism is used by control plane components (
kube-schedulerandkube-controller-manager) to acquire locks for leader election.
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:
etcd_mvcc_db_total_size_in_bytes(total file size)etcd_mvcc_db_total_size_in_use_in_bytes(space occupied by data in use)
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:
- etcd should use high-performance, low-latency local SSDs.
- Avoid sharing disks with other services that have high I/O loads, such as databases or log collectors.
- Regularly review
etcd_server_proposals_failed_total(the number of failed proposals; an unusual increase typically indicates frequent elections or disk latency) and request latency.
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:
- Stop all kube-apiserver instances to prevent inconsistent data from being written during restoration.
- Use
etcdutl snapshot restoreon every etcd node to rebuild a fresh data directory. - Start all etcd nodes and verify cluster health.
- 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 mechanism | Essential idea | Kubernetes mapping and application |
|---|---|---|
| Raft consensus | Writes require persistence by a majority; reads are linearizable by default | Use an odd number of members, three or five; losing quorum freezes the control plane while existing Pods continue running |
| MVCC + Revision | A globally increasing logical clock for version history | Object versions through resourceVersion; optimistic locking and conflict detection through 409 Conflict |
| Watch | Subscribe from a specified Revision and replay a stream of incremental events | Informer List-and-Watch and the kube-apiserver Watch Cache |
| Lease | A declaration of liveness with TTL and KeepAlive | Node 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.