Contents

Architecture & System Design › Distributed Systems

Distributed System

Many machines cooperating over an unreliable network.

Also known as: distributed systems, distributed computing, networked system, distributed application

A distributed system is a collection of independent machines that cooperate over a network to appear as one system. Almost every real backend today is one: a web app, a database, a cache, a queue and a few services already count.

Why build them

  • Scale: more machines handle more load and data than one can.
  • Availability and fault tolerance: if one machine fails, others continue.
  • Latency: put data and compute near users.
  • Organization: different teams own different services.

Why they’re hard

The difficulty comes from properties that a single machine doesn’t have:

  • Partial failure. Parts fail independently. A request might succeed on one machine and fail on another, and you can’t always tell whether a failed or slow call actually did its work. A timeout means “I don’t know”.
  • An unreliable network. Messages get lost, delayed, duplicated and reordered. Latency varies. The network can split the system into groups that can’t talk to each other (network partitions, CAP theorem). See the fallacies of distributed computing.
  • No shared clock and no global state. Machine clocks drift, so you can’t reliably order events by timestamp (clock skew, logical clocks).
  • Concurrency everywhere, with races across machines.
  • Agreement is expensive. Getting machines to agree on a value (a leader, an order of operations) needs consensus algorithms, which cost latency and are tricky to get right (consensus, Raft, Paxos).
  • Observability is harder. A request touches many machines, and the story is spread across logs and traces (distributed tracing).

Core techniques

ProblemTypical tools
Surviving failuresRedundancy, replication, failover (database replication)
Unknown outcomes of callsIdempotency, retries with backoff, timeouts (idempotency key, retry with backoff)
Cascading failureCircuit breakers, bulkheads, load shedding (blast radius)
Consistency across nodesConsistency models, quorums, consensus (consistency models, quorum)
Multi-step operationsSagas, outbox (saga pattern, transactional outbox)
CoordinationLeader election, locks and leases, fencing tokens (leader election, leases)
Finding other servicesService discovery (service discovery)

Practical principles

  • Assume failure at every call and every node. Design what happens when each dependency is slow, wrong or down.
  • Make operations idempotent.
  • Prefer simplicity. Every distributed component adds failure modes. A single database often beats a clever distributed design for a long time (monolith vs microservices).
  • Use proven building blocks (managed databases, queues, consensus services) instead of inventing coordination protocols.
  • Test with faults: kill nodes, delay and drop messages, partition the network (chaos engineering).
  • Measure and trace.

If you take one thing away: a remote call is not a local call. It can fail, be slow, be duplicated, or have an unknown result.