Chubby is the lock service Google built in the early 2000s so that systems such as GFS and Bigtable could elect a single master and store a few small pieces of critical metadata. Mike Burrows described it in the OSDI 2006 paper The Chubby lock service for loosely-coupled distributed systems, and that paper is worth studying for two reasons. It is one of the clearest accounts of how to wrap consensus in something application programmers will actually use, and it is unusually honest about how the system was misused once it was deployed.
This article explains Chubby from the paper; every number quoted is the paper's, from its time. It covers the cell and master, locks and sequencers, the consistent client cache, sessions, failover, scaling and what went wrong in practice, then compares it with ZooKeeper and etcd, which inherited most of these ideas, and a checklist for anyone building on a coordination service today.
Why a lock service and not a Paxos library
Google already had a Paxos library, so why a separate service? Systems start as prototypes with no thought for availability, and adding a lock call is far easier than restructuring a program as a replicated state machine. An elected primary must also advertise itself, and a lock service that stores small files does both. And consensus needs a quorum of the client's own replicas, typically three or five machines, whereas with a lock service a single client can take a lock and make progress safely; in the paper's phrase, the service acts as a generic electorate.
Two design choices follow. Chubby provides coarse-grained locks, held for hours or days, such as which server is the Bigtable master. Fine-grained, millisecond locks would let the service's load and failovers dominate its clients. And it favours availability and reliability over throughput.
The cell and the master lease
A Chubby cell is a small set of replicas, usually five, placed to fail independently; three must be up for the cell to work. The replicas run a consensus protocol to elect a master. A master must win votes from a majority and also a promise that those replicas will not elect a different master for a few seconds, the master lease. The lease is renewed as long as the master keeps winning a majority.
Writes go through consensus and are acknowledged at a majority. Reads are answered by the master alone, which is safe because while its lease is unexpired no other master can exist. Many Raft systems later adopted this lease-read optimisation; see leases for the general pattern and its clock assumptions.
Clients find the master by asking any replica listed in DNS. After a master failure, elections typically take a few seconds, though the paper saw up to 30.
Namespace, files and handles
Chubby looks like a simple file system. Names have the form /ls/foo/wombat/pouch: ls stands for lock service, foo is the cell, and the rest is interpreted inside the cell. Nodes are files or directories, permanent or ephemeral; an ephemeral node is deleted when no client has it open, which makes it a natural liveness indicator. Access control lists are themselves files, named in each node's metadata.
Files are small and read and written whole, deliberately, to discourage large files. Generation numbers let a client do compare-and-swap by passing the content generation it read to SetContents(). A handle from Open() names an instance of a file, not a path: if the file is deleted and recreated, old handles fail.
Handles can subscribe to events such as contents modified, child added or removed, master failed over and lock acquired. Events arrive after the change, so a read after a content-modified event returns the new data or newer.
Advisory locks, sequencers and lock-delay
Any node can be used as a reader-writer lock, exclusive for one holder or shared for many. Locks are advisory: holding one does not stop anyone touching the file, since the protected resources usually live in other services anyway.
The hard problem is the delayed request. Client X holds lock L and sends request R to a storage server, then X pauses or is partitioned and loses L. Client Y acquires L and does its work. Then R arrives and is executed as if L still protected it. Chubby's answer is the sequencer: an opaque byte string the holder obtains with GetSequencer(), containing the lock's name, the mode it was acquired in and the lock generation number. The holder sends it with each request, and the server rejects it if CheckSequencer() fails or a newer generation has been seen. This is the idea now generally called a fencing token.
For servers that cannot check sequencers there is a cruder guard: if a holder fails rather than releasing, the master withholds the lock for a lock-delay the client chose, up to one minute, so delayed requests can drain.
# Primary election on top of a Chubby-style service, with fencing.
def run_candidate(chubby, my_address):
h = chubby.open("/ls/cell/myservice/primary", mode="rw", events={"lock_acquired"})
h.acquire() # blocks until we are primary
h.set_contents(my_address) # advertise ourselves in the same file
seq = h.get_sequencer() # name + mode + generation
while session_is_safe(chubby):
for req in incoming_work():
storage.write(req.key, req.value, sequencer=seq) # storage checks it
# session in jeopardy or expired: stop sending work, do not assume we are primary
# On the storage server:
def write(key, value, sequencer):
if not chubby.check_sequencer(sequencer): # or compare with newest seen generation
raise StaleLeader()
apply(key, value)
Consistent client caching
With one master per cell and tens of thousands of clients, the client library's cache is what makes Chubby work. Clients cache file data and metadata, including the absence of files, as well as open handles and even locks. The cache is write-through and kept strictly consistent by invalidation.
When a file is to be changed, the master blocks the modification while it sends invalidations to every client that may have the data cached. The invalidations travel on KeepAlive replies, described next. The write proceeds once every client has acknowledged or let its cache lease expire. Meanwhile the node is treated as uncachable, so one round suffices and reads are never delayed. Invalidation beat update because updates can flow forever to clients that no longer care, and weaker consistency was rejected as harder for programmers.
Sessions, KeepAlives, jeopardy and grace
A session is the relationship between a client and a cell, kept alive by KeepAlive RPCs. Each session has a lease, a time interval during which the master promises not to end it unilaterally; the master may extend it but never shorten it. The trick is that the master holds each KeepAlive request blocked until the client's lease is nearly over, then returns it with a new lease timeout, by default extended by 12 seconds. The client immediately sends the next one, so there is almost always one KeepAlive waiting at the master. An overloaded master can extend leases further, up to around 60 seconds, to reduce KeepAlive traffic. The master returns a KeepAlive early to deliver events or invalidations, so a client cannot keep its session without acknowledging them.
The client keeps a more conservative estimate of the lease, allowing for transit time and bounded clock-rate differences. When that local lease expires, it empties and disables its cache and declares the session in jeopardy, telling the application through an event. It then waits a grace period, 45 seconds by default. If a KeepAlive succeeds in that time, it sends a safe event and re-enables the cache; otherwise the session has expired and API calls fail.
# Client-side session state machine (simplified)
state = "SAFE"
while True:
reply = keepalive(timeout=local_lease_remaining())
if reply.ok:
local_lease = conservative(reply.lease_timeout)
apply_invalidations(reply.invalidations); deliver(reply.events)
if state == "JEOPARDY":
state = "SAFE"; enable_cache(); emit("safe")
continue
if state == "SAFE" and local_lease_expired():
state = "JEOPARDY"; disable_and_flush_cache(); emit("jeopardy")
grace_deadline = now() + GRACE # 45 s by default
if state == "JEOPARDY" and now() > grace_deadline:
emit("expired"); fail_all_handles(); break
find_master_and_retry()One more property matters for correctness: once an operation on a handle fails because the session expired, every later operation on it fails the same way. Outages therefore lose only a suffix of a sequence of operations, never a random subset, so an application can make a multi-step change and mark it committed with a final write.
Master failover, step by step
When a master dies, it loses its in-memory state about sessions, handles and locks. The session timer is effectively stopped until a new master exists, which is equivalent to extending every client's lease. The paper lists what a new master does:
- Picks a new client epoch number that every call must carry, and rejects calls with older epochs, so it never acts on a delayed packet meant for a predecessor.
- Answers master-location requests, but does not yet process session operations.
- Builds in-memory structures for sessions and locks recorded in the database, extending session leases to the maximum the previous master may have granted.
- Allows KeepAlives, but no other session operations.
- Sends a fail-over event to every session, so clients flush caches (they may have missed invalidations) and applications learn that events may have been lost.
- Waits until each session acknowledges the event or expires.
- Allows all operations.
- Recreates handles created before the failover when clients use them, and remembers handles closed in this epoch so a delayed packet cannot recreate them.
- After an interval (about a minute), deletes ephemeral files with no open handles.
A worked timeline with the defaults: the master dies at t = 0, just after granting a client a lease ending at t = 8. The client's conservative estimate ends at t = 7; it enters jeopardy, flushes its cache and starts its 45-second grace period. A new master is elected at t = 15. The client's first KeepAlive to it is rejected for carrying the old epoch; its retry succeeds, the client receives a fail-over event and a new lease, and it sends a safe event to the application at about t = 16. The application saw a 9-second pause in Chubby calls and nothing else. The paper admits this rarely exercised code was a rich source of bugs.
Storage and scaling
Replicated Berkeley DB was replaced by a simple purpose-built database with a write-ahead log and snapshots, replicated by the same consensus protocol. Snapshots go to GFS in another building every few hours.
Clients are processes, not machines, and the paper saw 90,000 of them talking to one master. Since KeepAlives dominate traffic, scaling means talking to the master less: one cell per data centre, longer leases under load, aggressive caching, and protocol-conversion servers such as a DNS front end. Proxies and namespace partitioning by hash(D) mod N were designed but not yet in production.
How it was really used, and what went wrong
The biggest surprise was that Chubby became Google's main internal name service, used far more for naming than for locks. Short DNS TTLs overload DNS servers when thousands of processes look each other up; Chubby's invalidated cache stays fresh without polling. Bigtable used it to ensure a single master, store the bootstrap location of its data, discover tablet servers through ephemeral files, and hold schema and ACLs.
- Developers rarely consider availability. One system of hundreds of machines started recovery taking tens of minutes whenever Chubby elected a new master, multiplying a short blip a hundredfold. Many applications crashed on the fail-over event, which was only meant as a hint to rescan.
- No quotas. Teams stored growing amounts of data in Chubby. Google introduced a 256 KB file size limit and steered heavy users elsewhere.
- Abusive clients. Google added caching of file absence, artificial delays and design reviews before teams could use a shared cell.
- Transport mattered. TCP's back-off ignored lease deadlines and caused lost sessions during congestion, so KeepAlives moved to UDP.
Chubby, ZooKeeper and etcd
| Aspect | Chubby | ZooKeeper | etcd |
|---|---|---|---|
| Primitive | Locks plus small files | Wait-free znodes; locks are client recipes | Key-value with leases and transactions |
| Consensus | Paxos-based, master lease | Zab | Raft |
| Reads | Master only, plus consistent client cache | Any server; may be stale unless synced | Linearizable by default via leader |
| Notification | Invalidations and events on KeepAlive | One-shot or persistent watches | Watch streams with revisions |
| Fencing | Sequencer with lock generation | zxid or znode version as a token | Revision or lease id as a token |
| Availability | Proprietary to Google | Open source | Open source |
The resemblance is strong: sessions, ephemeral nodes, notifications and fencing tokens. The main fork is consistency of reads: Chubby paid for strict cache consistency with invalidation traffic, while ZooKeeper chose to scale reads across servers and let clients ask for freshness when they need it. The consensus foundations are covered in Paxos.
What to do next
- Read the 2006 paper itself, especially failover and the lessons.
- For every lock you take on ZooKeeper or etcd today, find where the protected resource checks a fencing token; if nothing checks, you have Chubby's delayed-request problem without its sequencer.
- Handle the coordination service's equivalents of jeopardy and expiry separately: pause work on disconnection, give up leadership only on expiry, and never crash on a reconnect event.
- Measure your coordination service's failover time under load and make sure your clients' grace or session timeout covers it.
- Set size limits and per-client quotas on your coordination store before teams start using it as a database.
- Review new users of a shared cluster before they go live, as Google learned to do, checking for polling, fan-out and large values.