The short answer
Quick answer: On a single computer, things either work or the whole machine fails. In a distributed system, part of it can fail while the rest keeps running, and you often cannot tell which. Messages are lost, delayed, duplicated or reordered. A remote machine that stops answering might be dead, slow or just unreachable, and these look identical. Clocks on different machines disagree. There is no shared memory and no global "now". Every guarantee you take for granted on one machine has to be rebuilt with timeouts, retries, replication and consensus, each of which introduces its own difficulties.
Why build them at all
Despite the pain, we distribute systems for good reasons:
- Scale. One machine cannot hold or serve everything.
- Reliability. One machine is a single point of failure.
- Latency. Users are spread across the world.
- Organisation. Many teams need to work on separate parts.
As soon as your application has a database on another host, a cache, or a third-party API, you have a distributed system.
The eight fallacies
In the 1990s, engineers at Sun Microsystems listed the false assumptions newcomers make, now known as the fallacies of distributed computing:
- The network is reliable.
- Latency is zero.
- Bandwidth is infinite.
- The network is secure.
- Topology doesn't change.
- There is one administrator.
- Transport cost is zero.
- The network is homogeneous.
Each one is something that is effectively true inside a single process and false across machines. Most distributed bugs trace back to code that quietly assumes one of them.
Problem 1: Partial failure
When a function call inside your program fails, you get an exception. When a call across a network gets no response, any of these may have happened:
- The request was lost.
- The remote server crashed before doing the work.
- It did the work and crashed afterwards.
- It did the work and the reply was lost.
- It is just slow, and the reply is on its way.
You cannot distinguish them. This is the essence of the Two Generals' Problem, and it is why a timeout tells you nothing about whether the operation happened.
Problem 2: The network
Networks drop packets, delay them unpredictably, deliver them twice and deliver them out of order. Cables are cut, switches are misconfigured, and cloud networks have brief outages.
A partition, where two groups of machines cannot reach each other while both keep running, forces the choice described by the CAP theorem.
There is no perfect timeout. Too short, and you declare healthy nodes dead and trigger needless failovers. Too long, and users wait while a truly dead node is not replaced.
Problem 3: Clocks
Each machine has its own clock, and they drift apart. Synchronisation helps but leaves errors of milliseconds or more, and a clock can even jump backwards when corrected.
So you cannot reliably order events on different machines by timestamp. A "last write wins" rule based on wall-clock time can silently discard the newer write. See why clocks can't be trusted in distributed systems.
Problem 4: No shared state
In one process, all threads see the same memory. Across machines, each node has only its own view, built from messages that may be late.
"What is the current value?" has no single answer. Different nodes may briefly disagree, and the system has to define what that means for users. That is the topic of eventual consistency.
Problem 5: Concurrency without a referee
Two users update the same record through two different servers at the same moment. On one machine, a lock or a database transaction would sort it out. Across machines there is no shared lock. Building one that stays correct under failures is surprisingly difficult; see how distributed locks work.
Problem 6: Agreement
Many tasks require nodes to agree on one thing: who the leader is, whether a transaction committed, what order operations happened in. This is consensus.
A famous result (the FLP impossibility result) shows that in a fully asynchronous network, no algorithm can guarantee consensus if even one node may crash. Practical algorithms such as Paxos and Raft get around it by using timeouts and requiring a majority of nodes, accepting that progress can stall while the network misbehaves.
Getting it wrong causes split brain: two nodes both believe they are the leader, both accept writes, and the data diverges.
Problem 7: Process pauses
A process can freeze for seconds without knowing it: a garbage collection pause, a virtual machine being moved, heavy swapping. When it resumes, it carries on as if no time has passed. It may still think it holds a lock or is the leader, when the rest of the system moved on long ago.
Problem 8: Failure is normal at scale
If one server fails once every three years on average, a fleet of a thousand servers sees about one failure per day. With thousands of disks, one is always failing. At scale, rare events happen constantly, so systems must be designed to keep running through failures, not to avoid them.
Failures can also spread:
- Retry storms. A slow service causes clients to retry, tripling its load and finishing it off.
- Cascading failure. One overloaded node fails, its traffic moves to the others, and they fail in turn.
There is also a more hostile category, Byzantine faults, where nodes send wrong or contradictory information because of bugs, corruption or malice. Most internal systems assume nodes are honest; blockchains and aerospace systems cannot.
Problem 9: Testing and debugging
Distributed bugs depend on timing and on rare failures occurring in a particular order. They may appear once in a million runs and vanish when you add logging. There is no single debugger and no single clock to order the logs by.
Tools that help include distributed tracing (see logs, metrics and tracing), fault injection and chaos testing, deterministic simulation, and formal specification. The Jepsen project tests databases under induced network failures and has found violations in many well-known systems.
How engineers cope
| Technique | What it addresses |
|---|---|
| Timeouts on every remote call | Unbounded waiting |
| Retries with exponential backoff and jitter | Transient failures, without retry storms |
| Idempotency | Makes retries safe |
| Replication | Node and disk failure |
| Quorums and consensus | Agreement without a single point of failure |
| Circuit breakers and load shedding | Cascading failure |
| Fencing tokens | Stale leaders and lock holders |
| Logical clocks | Ordering events without trusting wall clocks |
| Observability | Finding out what happened |
The practical advice
- Do not distribute until you must. One large server and a simple design avoid most of these problems.
- Use proven building blocks. Do not write your own consensus or replication.
- Assume every remote call can fail, hang or happen twice.
- Keep state in as few places as possible.
- Design for degradation. Decide what the system should do when a dependency is unavailable.
Frequently asked questions
What is a partial failure?
A situation where some components of a system fail while others keep working, and the working parts cannot be sure what state the failed parts are in.
What are the fallacies of distributed computing?
Eight common false assumptions, such as "the network is reliable" and "latency is zero", that lead to fragile distributed software.
Why can't you tell whether a remote server has crashed?
Because a crashed server, a slow server and a broken network all produce the same symptom: no response.
Are microservices a distributed system?
Yes. Every call between services is a network call and inherits all of these problems. See monolith vs microservices.
Conclusion
Distributed systems are hard because the comfortable assumptions of a single machine no longer hold: calls can fail ambiguously, time is not shared, and failure is constant. None of this can be eliminated. It can only be managed, with timeouts, retries, idempotency, replication and consensus, and with a healthy reluctance to distribute anything that does not need to be.
Related articles
- The Two Generals Problem: Why Perfect Communication Is Impossible
- The CAP Theorem Explained With Real Examples
- Why Clocks Can't Be Trusted in Distributed Systems
- How Consensus Algorithms Like Raft Keep Servers in Agreement
