DEV Community

Anusha Mukka
Anusha Mukka

Posted on

Consensus Is the Most Expensive Word in Your Architecture

Every write pays the quorum tax. Know exactly what you are buying before you embed Raft, etcd, or ZooKeeper.

Two of your three etcd nodes go dark during a network partition. Your application servers are all healthy. CPU is fine, memory is fine, the database is fine. But no leader can be elected, no lock can be acquired, no lease can be renewed, and every deploy gated on a lock stops dead. The three smallest processes in your fleet just vetoed the whole system, and they did it by design.

That scenario is a composite, but every element of it has paged someone. That is the consensus tax. It is the most expensive sentence in your architecture, and you pay it on every write.

Let me give you an example of the two bad options consensus replaces. Option one: a single coordinator that everyone asks. Simple to reason about, and one crash turns every agreement into a guess. Option two: every node decides for itself and you reconcile later. Fast, and the reconciliation is where the double charges, the duplicate leaders, and the 3 a.m. pages live. Consensus is the third option: get a majority of nodes to agree on the order of operations before any of them acts.

Here is the part every explainer skips. Consensus sells you exactly one thing: an agreed-upon sequence of operations. Not locking, not leader election, not safety in general. A log everyone agrees on. The price is a quorum round trip on every write, plus write unavailability whenever a quorum is unreachable. Most teams I have seen buy consensus for problems that only needed a lease, a compare-and-swap, or a fencing token, and then they pay the tax on operations that never needed the guarantee.

Consensus sells you exactly one thing: an agreed-upon sequence of operations. Everything else is machinery around that.

Learn What a Quorum Actually Guarantees

Start with the mechanism, because the design choices are not accidents. A consensus protocol like Raft or Paxos gives you two properties. Safety: at most one value is ever decided for each slot in the log, and once decided, it is never contradicted. Liveness: the system keeps making progress as long as a majority of nodes can talk to each other.

The safety property rests on one piece of arithmetic: any two majorities overlap. In a five-node cluster, a majority is three. Any two groups of three share at least one node. That shared node carries the memory of the earlier decision into the later round, so a new leader can never un-decide what was already committed. Leader election, log replication, snapshots, all of it is machinery to protect that one overlap.

And here is the honest caveat the textbooks bury in chapter one. Fischer, Lynch, and Paterson proved in 1985 that in a fully asynchronous network, no deterministic protocol can guarantee consensus if even one process can fail. That is the FLP result, and it is not a curiosity. It means Raft and Paxos do not guarantee termination. They terminate in practice because real networks deliver messages within bounded time most of the time. Your consensus cluster is standing on the assumption that your network is usually well-behaved. It usually is. Plan for the days it is not.

Run the Quorum Tax Yourself

Enough theory. Spin up three etcd nodes and watch the tax get collected in real time. This takes about five minutes with Docker.

docker network create etcd-net

for i in 1 2 3; do
docker run -d --name etcd$i --network etcd-net \
  -p 237$i:2379 -p 238$i:2380 \
  quay.io/coreos/etcd:v3.5.12 \
  /usr/local/bin/etcd \
  --name etcd$i \
  --listen-client-urls http://0.0.0.0:2379 \
  --advertise-client-urls http://etcd$i:2379 \
  --listen-peer-urls http://0.0.0.0:2380 \
  --initial-advertise-peer-urls http://etcd$i:2380 \
  --initial-cluster etcd1=http://etcd1:2380,etcd2=http://etcd2:2380,etcd3=http://etcd3:2380 \
  --initial-cluster-token my-etcd-cluster \
  --initial-cluster-state new
done

ENDPOINTS="http://localhost:2371,http://localhost:2372,http://localhost:2373"
etcdctl --endpoints=$ENDPOINTS endpoint health
etcdctl --endpoints=$ENDPOINTS put /config/feature-x on
Enter fullscreen mode Exit fullscreen mode

A few things are worth noting about this example. First, the three nodes find each other through the initial-cluster list and elect a leader; you did not configure one. Second, every write you issue goes to the leader, gets appended to its log, replicates to a majority, and only then commits. That replication round trip is the tax, and you pay it on every single write. (Any recent 3.5.x etcd image works here; the commands are the same.)

Now collect the tax. Kill two nodes and write again:

docker stop etcd2 etcd3
etcdctl --endpoints=$ENDPOINTS put /config/feature-x off
# Error: etcdserver: request timed out
etcdctl --endpoints=$ENDPOINTS --consistency=s get /config/feature-x
# /config/feature-x
# on
Enter fullscreen mode Exit fullscreen mode

Notice what just happened. The surviving node is perfectly healthy. The data is intact. And it refused your write anyway, because one node out of three is not a majority, and etcd would rather say nothing than risk two halves of the world disagreeing. That refusal is the safety property doing its job. It is also your deploy pipeline, your leader election, and your lock service all going dark at once. The stale read with --consistency=s still works, which is the consolation prize: you can read old data from a minority, but you cannot agree on anything new.

When you are done, clean up with docker rm -f etcd1 etcd2 etcd3. You just watched a healthy node choose silence over disagreement. Every consensus system you will ever operate makes exactly that trade.

The Overlap Is the Whole Trick

If the Docker demo is the tax, this Python is the receipt. The entire safety argument for majority quorums fits in six lines of standard library:

from itertools import combinations

nodes = {"n1", "n2", "n3", "n4", "n5"}
majority = len(nodes) // 2 + 1

quorums = [set(q) for q in combinations(nodes, majority)]
assert all(a & b for a, b in combinations(quorums, 2))
print(f"All {len(quorums)} possible majorities pairwise overlap.")
Enter fullscreen mode Exit fullscreen mode

Run it. It passes. Any two majorities share at least one node, so a decision committed by quorum A is always visible to any later quorum B. That is why the number is "more than half" and not some other fraction: it is the smallest group with the overlap property, which means it is also the largest group that can still form when nodes fail. The quorum size is not a tuning parameter. It is the answer to a math problem.

Now the cheaper alternative most of your problems actually need. If all you want is an atomic update to one key, you do not need a log everyone agrees on. You need compare-and-swap:

def increment_counter(store, key):
    while True:
        value, version = store.get(key)       # read value + version
        if store.compare_and_swap(key, value + 1, version):
            return value + 1                  # we won the race
        # someone else wrote first; reread and retry
Enter fullscreen mode Exit fullscreen mode

This is a single-key atomic update with no leader, no log, no quorum. etcd gives it to you as a transaction on one key, Consul as Check-And-Set, DynamoDB as a condition expression, Postgres as UPDATE ... WHERE version = .... The walkthrough matters here: notice that the retry loop is the whole concurrency control. No consensus protocol was invoked, because there was no ordering problem to solve, only a last-writer question about one key. Reach for consensus when you need many operations in an agreed order. Reach for compare-and-swap when you need one key to change atomically.

Draw the Decision Before You Draw the Architecture

When a design review reaches for consensus, sketch this flow first:

Do readers need to agree on the ORDER of writes?
|
+-- YES --> consensus: one etcd cluster, a Raft library,
|           or a managed control plane. Never hand-rolled.
|
+-- NO --> Do you need at-most-one holder of a thing?
           |
           +-- YES --> lease + fencing token.
           |           The lease says who holds it now;
           |           the fencing token stops the deposed
           |           holder from acting.
           |
           +-- NO --> Is it an atomic update of one key?
                      |
                      +-- YES --> compare-and-swap.
                      |
                      +-- NO --> gossip, CRDTs, or plain
                                 replication. Converge, don't agree.
Enter fullscreen mode Exit fullscreen mode

And when you do pay for consensus, know the pipeline you are buying:

client --> leader --> append to local log
                     --> replicate to followers
                     --> majority ack --> COMMIT
                     --> apply to state machine --> reply
Enter fullscreen mode Exit fullscreen mode

Every arrow between "append" and "COMMIT" is a network round trip you pay per write. That is the tax, drawn honestly. If your write path goes through this pipeline, your write latency has a floor of one quorum round trip, and your write availability has a ceiling of quorum reachability. Both are true on day one and on day one thousand.

Where This Breaks

Unhedged, because this is the section that saves you a page:

  • FLP means termination is never guaranteed, only likely. A network that stops delivering messages on time can stall a Raft cluster indefinitely. Your uptime math quietly assumes a well-behaved network.
  • Lose quorum, lose writes. Two of three nodes unreachable and the cluster becomes a read-only museum of old data. There is no degraded mode for agreement.
  • The control plane becomes the blast radius. etcd does not just hold your config; your locks, leader elections, and service discovery lean on it. When it wedges, the wedge is systemic, which is exactly what the cold open described.
  • Write throughput tops out at what one leader can push through a quorum round trip. Consensus does not scale writes horizontally. If your write volume needs sharding, the consensus cluster is the wrong place for that data.
  • Consensus does not fence. A leader that gets partitioned away does not know it was deposed. Without fencing tokens checked on every write, the old leader keeps acting with stale authority. Agreement on order and protection from impostors are two different problems.
  • You have to operate it. The etcd maintainers will tell you the most common production failure is slow disks. Keep the store small, defrag on schedule, and monitor disk latency like it is a load-bearing wall, because it is.

Build It If, Skip It If

Build it if you need a linearizable source of truth for configuration, locks, or membership, and you are embedding an existing engine: etcd, Consul, ZooKeeper, or a tested Raft library. Build it if your team will operate exactly one consensus cluster and point every agreement problem at it, instead of growing a new one per service.

Skip it if you need at-most-once execution. That is idempotency keys, which I wrote about in Exactly-once delivery is a lie. Skip it if you need mutual exclusion across crashes. That is leases plus fencing tokens, from the clock-skew piece. Skip it if your replicas can converge instead of agreeing. That is gossip or CRDTs, and they ask nothing of your write path.

The minimal viable version fits in an afternoon: one three-node etcd cluster (or your cloud's managed lock service), application data kept strictly out of it, and a written rule that nobody ships their own agreement protocol. Steelman the hand-rolled impulse for a second: it feels lighter than operating etcd. It is lighter right up until the first split brain, at which point you have written a worse Raft with no tests and a production incident attached.

The cheapest consensus cluster is the one you never built because compare-and-swap was enough.

Point It at Your Own Write Path This Week

This week, find one place in your system that waits on consensus and ask what invariant it actually needs. Ordered writes, or just one agreed value? One agreed value, or just at-most-one actor? Most of the time you will find the tax is being collected on a road nobody drives. Move the agreement down to the cheapest mechanism that preserves the invariant, keep the consensus cluster small and boring, and spend the savings on the fencing tokens everyone forgets.

What is the smallest consensus cluster you have ever seen take down the largest system?

Resources

  1. In Search of an Understandable Consensus Algorithm (Ongaro and Ousterhout, USENIX ATC 2014): the Raft paper, genuinely readable.
  2. The Raft consensus algorithm site: includes the interactive visualization that makes leader election click.
  3. etcd documentation: the operations guide, especially the hardware and disaster-recovery sections.
  4. "Impossibility of Distributed Consensus with One Faulty Process" (Fischer, Lynch, and Paterson, 1985): the FLP result, short and worth the effort.
  5. Designing Data-Intensive Applications (Martin Kleppmann): chapters 8 and 9 are the best long-form treatment of consensus trade-offs.
  6. "The Part-Time Parliament" (Leslie Lamport): the original Paxos paper; read it after Raft, as history.

Top comments (0)