Single-leader consensus won the last fifteen years of distributed database design for a good reason: it is comprehensible. Raft in particular succeeded because a competent engineer can read the paper and reason correctly about the resulting system, which was never quite true of Paxos. That comprehensibility bought the industry a decade of reliable replicated state machines.
It also embedded an assumption into the API of nearly every distributed database built in that period — that somewhere there is a primary, and that writes go there. The research direction now reaching production questions whether that assumption needs to survive, and the consequences land less on database internals than on the shape of the APIs built above them.
Where the Leader Becomes the Constraint
The leader in a Raft group is a serialization point by design. Every write passes through one node, which orders it, replicates it to a quorum, and acknowledges. Within a single datacenter this costs almost nothing — the leader is a millisecond away from its followers, and having one authoritative sequencer simplifies everything downstream.
Stretch the same group across regions and the cost becomes structural. A client in Singapore writing to a group whose leader sits in Virginia pays a round trip to Virginia before the leader even begins replication, then waits for the leader’s quorum to respond. The client’s latency is determined not by its own distance to a quorum but by its distance to a specific node that happens to hold an election result. Two clients on the same continent can see very different write latencies purely because of where an election landed.
The standard mitigation is to shard aggressively and place leaders deliberately. Multi-Raft — the approach in TiKV, CockroachDB, and YugabyteDB — runs many independent Raft groups, one per data range, and distributes their leaders across the cluster. This is genuinely effective: it converts one global bottleneck into thousands of small ones and allows leader placement to follow access patterns. It does not eliminate the underlying property. Each individual key still has one leader at any moment, and a transaction spanning ranges whose leaders sit in different regions still pays for the geography.
Research on flexible quorums and WAN-oriented protocols — Flexible Paxos, WPaxos, EPaxos and its descendants — attacked the problem from the other direction by asking whether a distinguished leader is necessary at all. EPaxos established the key insight: commands that do not interfere with each other do not need to be ordered relative to each other, so a protocol can commit non-conflicting operations in one round trip from whichever replica received them. Order is imposed only where operations actually conflict.
Accord Puts It in Production
That line of research has now shipped. Apache Cassandra’s Accord protocol, proposed as CEP-15 and delivered in Cassandra 6, provides general-purpose distributed transactions with no elected leaders. It offers strict serializable isolation across multiple partitions and tables in a single round trip, and preserves the availability characteristics Cassandra is designed for under minority failures.
The engineering timeline is worth noting for anyone assessing maturity: development ran roughly three and a half years before the team began validating against real workloads in early 2025, and the remaining work at that point concerned edge cases — hybrid logical clock reuse after restart, ByteOrderPartitioner support, and live migration between Paxos and Accord — rather than the core protocol. This is a long-gestation feature that arrived closer to finished than most.
Cassandra is an instructive place for this to land. Its earlier transaction support used Paxos for lightweight transactions, with the performance profile that implies, and its default replication model was explicitly non-consensus peer-to-peer — the same lineage as DynamoDB and Couchbase, where availability came at the cost of not offering strong multi-key guarantees at all. Accord closes that gap without reintroducing a leader.
The API Consequence
Here is where this matters for anyone building above a database rather than inside one.
Single-leader systems leak their topology into the API. If writes must reach a leader, then either the client library must know where the leader is, or a proxy must route on the client’s behalf, or the application must accept unpredictable write latency. In practice this produces API concepts that exist purely to describe consensus mechanics: primary regions, write endpoints distinct from read endpoints, follower-read modes with explicit staleness bounds, and documentation explaining which operations are safe against a replica. Every one of those is a database implementation detail promoted into the public contract because it could not be hidden.
Leaderless commit changes what can be hidden. When any replica can coordinate a transaction and non-conflicting operations commit in one round trip, write latency becomes a function of the client’s distance to a quorum — a property of geography, which application architects can reason about and plan capacity around — rather than a function of where an election landed, which they cannot. A single logical endpoint becomes an honest abstraction rather than a proxy hiding a redirect.
The caution is that this is a change in the shape of the latency, not a removal of it. A quorum still has to be reached, and conflicting operations still require ordering, which means contended keys still serialize. Workloads with heavy write contention on a small key range will not find that leaderless consensus has repealed the speed of light or the cost of coordination. What changes is that uncontended cross-region writes stop paying a leader tax they never needed to pay, and the API stops needing a vocabulary for explaining that tax to users.
For teams designing cloud data APIs now, the practical move is to stop treating “which region is primary” as a permanent fixture of the interface. It has been a necessary disclosure for a decade because the underlying protocols made it unavoidable. That is ceasing to be true, and an API that hard-codes the concept will be describing an implementation constraint its own storage layer no longer has.
