computer science//distributed systems

A distributed system is a set of computers that cooperate by exchanging messages over a network that delays, loses and reorders them, and that must still behave toward its users as a single machine. Behind any serious robot fleet there is one: the database of orders, tasks and maps is replicated on three or five servers so that losing one does not stop the warehouse, and the Kubernetes cluster that runs the fleet software keeps its own state in a replicated store as well.


A distributed system is a set of computers that cooperate by exchanging messages over a network that delays, loses and reorders them, and that must still behave toward its users as a single machine. Behind any serious robot fleet there is one: the database of orders, tasks and maps is replicated on three or five servers so that losing one does not stop the warehouse, and the Kubernetes cluster that runs the fleet software keeps its own state in a replicated store as well.

What makes the subject hard is what one computer gives for free. There is no shared memory and no shared clock, so a server cannot tell a peer that crashed from a peer that is slow or a link that dropped; a message may arrive twice, late or never; and two servers can each believe they hold the latest value. The field is the set of answers to one question, what can still be guaranteed when any message may be late and any node may stop. Its word consensus means exact agreement by vote, and it shares only the name with the averaging consensus protocol of multi-agent control.

Distribute only when one server is not enough.

A single server with good backups covers many cases; replication buys availability, scale or geography, and pays for them with failure modes a single machine never has. The same argument decides fleet architecture for robots.

Two impossibility results draw the boundary before any algorithm is chosen. FLP impossibility says that in a fully asynchronous network no deterministic algorithm can guarantee agreement if a single node may crash, and the CAP theorem says that during a network partition each side must choose between answering with possibly stale data and refusing to answer. Practical systems live inside these limits, usually by giving up progress for a while and never giving up safety.

The workhorse is state machine replication: every server applies the same commands in the same order, so all of them compute the same state. Agreeing on that order is distributed consensus, solved in production by Paxos and by Raft, where a majority keeps the cluster going (2f+12f+12f+1 servers survive fff crashes).

The failure model sets the price. Tolerating nodes that only stop costs 2f+12f+12f+1 replicas; tolerating nodes that lie or misbehave arbitrarily, Byzantine fault tolerance, costs at least 3f+13f+13f+1 and many more messages, which is why it appears where participants cannot trust each other (blockchain consensus) or where one faulty unit must not corrupt a vote (avionics).

Each layer of a fleet has its own clock rate: motor control at hundreds of hertz on a microcontroller, the fleet layer at 1 to 10 Hz on board, the warehouse planner at about one cycle per second on a local server, and the replicated state in a Raft cluster beneath them. Systems that blur those layers behave strangely in ways no single layer explains (hierarchical control).