What Are Distributed Systems and What Core Trade-offs Should You Understand?
A distributed system is a group of independent computers that coordinate over a network to act, from the user's point of view, like a single system. You need this model when one machine can no longer handle the load, when you can't tolerate a single point of failure, or when users are spread across regions and need low latency. The trade-offs below are what make distributed systems hard: you gain scale and resilience, but you take on partial failure, coordination cost, and hard choices between consistency and availability.
Why distribute a system at all
Four forces usually drive the decision, and they often conflict:
- Scalability — spread work across many machines instead of buying one bigger one.
- Fault tolerance — if one node dies, the system keeps serving.
- Latency — put data and compute near users so requests travel shorter distances.
- Geographic reach — serve users in multiple regions under one logical system.
Each of these adds coordination. The more nodes you have, the more ways the system can partially break.
Core concepts you'll reason with
Replication
Keep copies of data on multiple nodes. Replication improves availability and read throughput, but every copy must eventually agree. The question is when — synchronously (slower, more consistent) or asynchronously (faster, risk of stale reads).
Partitioning (sharding)
Split data across nodes so no single node holds everything. Partitioning scales writes and storage, but cross-partition queries and transactions get expensive, and a hot partition can bottleneck the whole system.
Consistency
How fresh must a read be? Strong consistency means every reader sees the latest write; eventual consistency means replicas converge over time. Most large systems sit somewhere between these, choosing per operation.
Availability
The fraction of time the system answers requests. High availability usually means tolerating node and network failures without human intervention.
Consensus
Getting multiple nodes to agree on a single value or ordering — used for leader election, configuration, and replicated logs. Consensus protocols are the machinery behind "who is in charge" and "what order did things happen."
The CAP theorem, in practice
CAP says that during a network partition (nodes can't talk to each other), a system must choose between consistency and availability:
| Choice during a partition | Behavior | When it fits |
|---|---|---|
| Consistency (CP) | Reject or block requests rather than return stale data | Financial ledgers, inventory counts, anything where wrong data is worse than no data |
| Availability (AP) | Keep answering, accept that some reads may be stale | Feeds, catalogs, analytics where a slightly old answer is fine |
The key word is during a partition. When the network is healthy, you can often have both. Real systems frequently make this choice per operation rather than for the whole system — a checkout path may favor consistency while a recommendation path favors availability.
Why partial failure is the hard part
In a single machine, failure is usually total: it's up or it's down. In a distributed system, some nodes work while others don't, and the network itself can drop, delay, or reorder messages. That means:
- A request can time out even though the work succeeded — so retries must be safe (idempotent).
- You often can't tell the difference between a slow node and a dead one.
- Clocks across machines disagree, so "what happened first" is genuinely ambiguous.
This is why distributed systems demand careful design around retries, timeouts, and idempotency rather than assuming clean success or failure.
Grounding the concepts in real components
- Databases — replicated and partitioned stores are the canonical example; their consistency settings are where CAP shows up in your config.
- Caches — a distributed cache trades a bit of staleness for large latency and load wins.
- Microservices — splitting an app into services is a distributed system by construction, which is why service-to-service calls need the same failure handling.
If you want structured practice applying these ideas, System Designer offers fundamentals paths, whiteboards for sketching architectures, and guided project templates for full system design documentation — useful for turning the concepts above into working designs and interview-ready reasoning.