Shelf 1 · Computing Foundations · 6 / 45
Distributed Systems — History, Time, Failure, and Replication
From ARPANET to logical clocks, Byzantine faults, FLP, CAP, and replicated state machines: the foundations that long predate blockchain.
Check this article’s sources (11)Article brief
How do computers with no shared clock or manager keep one service alive while messages lag and individual machines fail?
A useful mental model
Imagine musicians in separate rooms trying to play one piece together, coordinating only through written messages that may arrive late. Time, ordering, and partial failure become visible at once.
Where the analogy stops
Distributed systems rely on explicit fault models and consistency rules, not human intuition, and different systems promise different things.
You will read a cloud outage and Bitcoin the same way: not by counting machines, but by asking what assumptions let them cooperate.
Open the glossaryArticle contents12 chaptersJump to a chapter
1What makes a system distributed
A distributed system is a set of separate processes that exchange messages and yet present one service or computation to the user. Leslie Lamport’s 1978 paper puts separated processes and non-negligible communication delay at the center of the definition. What decides the matter is not the number of machines: it is that no process can know another process’s state instantly.
Parallel work inside one large machine overlaps with distributed computing, but the two terms are not interchangeable. Across a network, a message can be delayed, lost, reordered, or delivered to only one side. One component can fail while the rest keeps running, and different machines can see the same events in different orders.
| Dimension | Parallel work on one machine | Distributed system |
|---|---|---|
| Communication | Shared memory or a fast interconnect | A network with delay, loss, and reordering |
| Time | Comparatively easy to coordinate | No perfectly shared clock can be assumed |
| Failure | The whole machine often fails together | One part can fail while the others stay live |
| Administration | One control boundary | May span organizations or anonymous participants |
2Keep the neighboring terms apart
Distributed, decentralized, peer-to-peer, grid, cloud, and blockchain name related but different properties. Distribution is about where computation and state live. Decentralization is about where control and authority live. A service can be spread across data centers in many countries and still be controlled by one company.
Peer-to-peer describes a network in which participants do not have to be permanently split into clients and servers. Grid computing pools resources across sites or organizations; cloud computing hands out provisioned resources on demand. A blockchain is a family of replicated, tamper-evident shared-ledger designs, and it occupies one region of the much larger field of distributed systems.
| Term | The main question it answers | What it does not guarantee |
|---|---|---|
| Distributed system | How do state and processing coordinate? | Decentralized control or open entry |
| Peer to peer | How are node roles and links arranged? | Consensus, a ledger, or anonymity |
| Volunteer computing | Whose compute can a project borrow? | Agreement among participants |
| Blockchain | How is shared history ordered and verified? | That every distributed workload needs one |
3From ARPANET to an open-participation ledger
ARPANET connected geographically separated computers over a packet network in 1969. Calling it “the first distributed system” overstates the case. A network provides communication; the research that followed had to work out how to coordinate state, replication, failure, and computation on top of that layer.
From the late 1970s through the 1990s, logical clocks, Byzantine faults, impossibility results, Paxos, and state-machine replication shaped the theory. In the 2000s, MapReduce spread bulk data processing across commodity clusters, and Dynamo documented a data store built around high availability. These systems made distribution practical for goals and trust models quite unlike Bitcoin’s.
The 2008 Bitcoin paper combined peer-to-peer communication, hash-linked history, proof of work, and incentives to address double spending in electronic cash. What was new was not distributed computing itself but the way existing pieces were put together to maintain a transaction history among open participants.
4There is no universal “now”
In a distributed system, asking which of two events happened first is not always a meaningful question. If an event at process A sends a message that arrives before an event at process B, the two are causally ordered. If no communication path connects them, they are concurrent, and a tiny difference in wall-clock readings creates no meaningful causal relationship.
Lamport’s happened-before relation and logical clocks number events in a way that is consistent with causality. The implication runs one way only: if `a → b`, then `C(a) < C(b)`. A smaller clock value on its own does not prove that one event caused the other. Vector clocks can keep more information about concurrency, at the cost of extra metadata.
This matters for double bookings, balances, file updates, and the arrival of blocks. In Bitcoin, two miners can find valid tips at almost the same moment, and different nodes can briefly see different heads. The protocol converges on an ordering through its chain-selection rules, not by assuming a perfect physical clock.
5Partial failure: broken, or just slow?
The difficulty peculiar to distribution is partial failure. When a node does not answer, an observer cannot tell for certain whether it crashed, the network partitioned, the reply was lost, or the node is merely slow. A timeout is a useful way to raise suspicion; it is not proof of death.
| Failure model | Observable behavior | Typical defenses |
|---|---|---|
| Crash | A process stops responding | Replication, leader change, re-execution |
| Omission | Some sends or receives disappear | Retries, deduplication, acknowledgments |
| Partition | Live groups cannot communicate | Quorums; choose consistency or response |
| Byzantine | Conflicting or arbitrary behavior | Authentication, redundant checking, BFT |
| Transient | State is temporarily corrupted | Self-stabilization and resynchronization |
A Byzantine fault does not require anyone acting maliciously. Software bugs, bad memory, corrupted traffic, and a compromised node can all produce arbitrary behavior outside the specification. The fault model you choose sets both the replication factor and the communication cost a protocol needs.
6Keep safety separate from liveness
Guarantees in a distributed system are usually split into safety and liveness. Safety means that something bad never happens: two replicas never commit different values at the same log position, and an invalid transaction is never accepted. Liveness means that something good eventually happens: a request completes, or the system finally decides.
One design may halt during severe network trouble in order to preserve safety. Another may keep answering and allow temporary inconsistency. “Never stops” and “stays correct” are not the same measure, and a system specification has to say which property it preserves under which fault.
Bitcoin draws the same line. Full nodes rejecting invalid blocks is a safety property. The continued production of valid blocks is a liveness property, and it rests on assumptions about network propagation and honest hash power.
7Replicated state machines: same order, same state
Replication is the basic tool for reliability, but copying data to several machines is not enough on its own. If concurrent updates are applied in different orders, the replicas drift apart. State-machine replication starts every deterministic replica from the same state and applies commands in the same order, so the replicas behave as one service.
The central job of a consensus algorithm is to decide which command occupies each position in the log. Paxos and Raft build replicated logs among known servers that may crash. PBFT extends the model to replicas that may behave arbitrarily. Before comparing algorithms, work out the membership, the fault bound, and the network assumptions.
Bitcoin full nodes also recompute UTXO state from shared validation rules, but the proposers are miners anyone can join, and history is selected by cumulative proof of work. That is a different membership model from a registered set of enterprise replicas.
8What FLP actually makes impossible
The 1985 Fischer–Lynch–Paterson result, known as FLP, shows that in a fully asynchronous message system no deterministic consensus protocol can guarantee termination in every admissible execution if even one process may crash. It does not say that distributed consensus is impossible in practice.
The theorem stacks up strong conditions: unbounded message and process delay, no reliable way to tell delay from failure, determinism, and termination in every execution. Working protocols make progress by adding an eventual timing assumption, randomization, failure detectors, or operational timeouts.
Dwork, Lynch, and Stockmeyer formalized partial synchrony, in which useful timing bounds exist but may only become known eventually. The lasting lesson is to read what a protocol assumes about time before reading its name.
9CAP is more precise than “pick two”
CAP says that during a network partition, a service cannot guarantee both atomic consistency, meaning behavior equivalent to one up-to-date copy, and availability, meaning that every request to a non-failing node receives a response. Partition tolerance is not a decorative feature to drop when the environment can partition.
The familiar triangle saying “choose two of consistency, availability, and partition tolerance” is memorable, but it can suggest a permanent choice made at design time, one that holds even while the network is healthy. Real systems decide at a finer grain: per operation, per data item, per latency budget, and per affected part of the network.
The C in CAP is not a vague promise that data will look similar eventually, and the A is not ordinary annual uptime. Keep the formal properties of the paper separate from a product’s service-level objectives.
10Clusters, grids, and clouds changed the scale
The 2004 MapReduce paper hid input partitioning, task placement, re-execution of failed tasks, and machine-to-machine communication behind a programming model for large commodity clusters. Treating failure as routine and recovering by re-running work mattered as much as raw parallel speed.
Amazon’s 2007 Dynamo paper documented a highly available key-value store built for always-on services such as shopping carts. Consistent hashing, vector clocks, sloppy quorums, and read repair expressed deliberate trade-offs rather than one universally correct consistency model.
Grid and volunteer computing took a different path, pooling unused resources scattered across organizations and homes. They hand out separable work and validate the results that come back, rather than building a shared transaction ledger. The Volunteer Computing article follows that lineage in depth.
11What Bitcoin inherited, and what it changed
Bitcoin inherited peer-to-peer messaging, replicated state, cryptographic hashing, digital signatures, and the old distributed problem of ordering events. Unlike a registered server cluster, though, it had to let anyone in and hold up when a single actor creates many pseudonymous identities.
Proof of work is not one node, one vote. It turns verifiable computational resources into weight for proposing history. Each node follows the valid chain that took the most cumulative work. Under the model assumptions, such as attacker hash power staying below the honest share, deeper confirmation lowers the probability of reversal; it does not create immediate mathematical finality.
Miners propose an ordering, but no amount of work makes an invalid block valid. Every full node checks signatures, double spending, issuance, and script rules for itself. Sybil resistance through scarce resources and validity through rule checking are two separate jobs.
12Six questions for reading any distributed design
Before you trust a product label or the word “decentralized,” ask six structural questions.
| Question | What to identify |
|---|---|
| Who participates? | Known servers, several organizations, or anonymous open entrants |
| What is shared? | Results, files, state, an operation log, or an asset ledger |
| What can fail? | Crash, partition, Byzantine behavior, or operator shutdown |
| What does time mean? | Synchronous, asynchronous, partial synchrony, deadlines |
| What counts as correct? | Safety, consistency, verifiability, finality |
| What sustains progress? | Quorums, a leader, retries, incentives, scarce resources |
BOINC, Raft, and Bitcoin are all distributed, and they solve different problems. BOINC sends independent jobs to untrusted home computers, runs a project-specific validator, and cross-checks replicated answers when it needs to. Raft orders a log among known servers. Bitcoin maintains a history of value transfers among open participants. Connecting the three without erasing those differences is the shortest route to understanding distributed computing.
Primary sources
- DARPA — ARPANET
- RFC Editor — RFC 1: Host Software
- Leslie Lamport — Time, Clocks, and the Ordering of Events in a Distributed System
- Lamport, Shostak, Pease — The Byzantine Generals Problem
- Fischer, Lynch, Paterson — Impossibility of Distributed Consensus with One Faulty Process
- Dwork, Lynch, Stockmeyer — Consensus in the Presence of Partial Synchrony
- Fred Schneider — Implementing Fault-Tolerant Services Using the State Machine Approach
- Gilbert and Lynch — Brewer’s Conjecture and the Feasibility of CAP Web Services
- Google Research — MapReduce: Simplified Data Processing on Large Clusters
- Amazon Science — Dynamo: Amazon’s Highly Available Key-value Store
- Bitcoin whitepaper
Read next
Volunteer Computing — The History of SETI@home and BOINC18 min readRelated topics
Go deeper
Citation
- Title
- Distributed Systems — History, Time, Failure, and Replication
- Source
- Bitcoin Library (bitcoin.ne.jp)
- Canonical URL
- https://bitcoin.ne.jp/en/learn/distributed-systems
- Author
- KK siiiiiixth
- Topic
- distributed-systems
- Published
- Updated
- Last verified
- Editorial policy
- https://bitcoin.ne.jp/en/editorial-policy
- About
- https://bitcoin.ne.jp/en/about
- License
- Content reuse terms
Operator-owned article text, original diagrams, and public data may be used for citation, summarization, indexing, search, RAG, machine analysis, and AI model training. When content is presented to readers, identify Bitcoin Library and the applicable canonical URL where technically practicable.