What is a distributed system?
A distributed system is a group of computers connected by a network that work together, so that users see one service. The computers pass messages to each other. The work (functions) and the information (data) are spread over several machines instead of sitting on one.
The client-server idea
The simplest shared design is client-server. A client (your phone or laptop) asks for something. A server is a computer that waits for requests and answers them: it sends a web page, an email or a file. One function (serving pages) is done by the server, another (showing the page) by the client. This is how web browsing, e-mail and online shopping work.
The problem with one server
A server can answer only so many requests per second. If more come, they wait in a queue and every user waits longer. If that single server breaks, the whole service stops. This is a single point of failure.
Sharing the work: many servers and a load balancer
To cope with more users we add servers. Now we need to decide which server gets each request. A load balancer is a traffic manager: it receives each request and sends it to a server that is not too busy. Common ways: take turns (round robin), or pick the server with the shortest queue.
This gives:
- Speed: shorter queues, so quicker answers.
- Scalability: if there are more users, add more servers.
- Different jobs on different machines: for example, one group of servers shows pages, another stores the data, another sends e-mails.
The load balancer checks which servers are alive. A broken server is skipped, so users are not sent to it.
Sharing the data: copies and pieces
Replication: keep copies
Replication means keeping copies of the same data on several machines, maybe in different cities. Benefits: the service keeps working if one machine fails (fault tolerance), and a user can read from a copy that is nearby, which is faster (lower latency).
Partitioning: keep pieces
If the data is too big for one machine, we cut it into pieces (for example by name A-F, G-M ...) and store each piece on a different machine. This is often called sharding. Each machine holds less, and each handles fewer requests.
The cost: keeping copies in step
If two copies are changed at different places, they may differ for a short time. Systems need rules to bring them back in step (consistency). Also more machines mean more network messages, more things that can fail and more to protect. So distribution is a trade-off, not a free gift.
Peer-to-peer and where you meet these ideas
Peer-to-peer (P2P)
In a peer-to-peer network there is no central server. Every computer (a peer) is both a client and a server. Peers swap pieces of a file with each other, so the more peers there are, the more helpers there are. File-sharing networks and some video calls work like this. If one peer leaves, the others carry on. A problem: it is harder to control who shares what, and to protect against harmful files.
Everyday examples
- The web: a browser (client) talks to web servers; popular sites use many servers and a load balancer.
- DNS: the "phone book" of the internet is spread over many servers around the world.
- Content delivery networks: copies of videos and images are kept near users.
- Cloud services: you rent computing and storage that is spread over many machines.
- E-mail and messaging: servers store and forward messages.
Key formulas and definitions
- Waiting requests after t seconds ≈ (arrival rate − service rate) × t, when arrival rate is bigger
- Capacity of n equal servers = n × (capacity of one server)
- Chance all n independent copies are down = (chance one is down)ⁿ
- Time to download ≈ file size ÷ total speed from all senders
- Key terms: client, server, peer, load balancer, replication, sharding, latency, fault tolerance
Worked examples
1. One server answers 1.2 requests per second. 6 requests arrive every second. How many requests are waiting after 10 seconds?
Waiting grows by 6 − 1.2 = 4.8 requests every second. After 10 s: 4.8 × 10 = 48 requests are waiting.
2. Three equal servers each answer 1.2 requests per second. What is their total capacity? If 3 requests arrive per second, how busy are they?
Capacity = 3 × 1.2 = 3.6 requests/s. Busy = 3 ÷ 3.6 ≈ 0.83, so about 83%. They can keep up, so queues stay short.
3. Each server is working 90% of the time, independently. One server is down 10% of the time. With 3 copies of the data on 3 servers, how often is at least one working?
All three down at once: 0.1 × 0.1 × 0.1 = 0.001 (0.1%). So at least one is up 1 − 0.001 = 0.999, which is 99.9% of the time. One server alone gave only 90%.
4. A table of 1 000 000 customer records is split equally over 4 servers (sharding). How many records does each server hold? What happens to a lookup of one customer?
1 000 000 ÷ 4 = 250 000 records each. A lookup goes only to the one server that holds that customer, so the other three are free for other work.
5. A 600 MB file is on one server that sends 8 MB/s. How long to download? Now 4 peers each send 8 MB/s at once. How long now?
One server: 600 ÷ 8 = 75 s. Four peers: total speed 4 × 8 = 32 MB/s, so 600 ÷ 32 = 18.75 s, about 19 s.
6. A signal travels about 200 000 km/s in fibre. A user is 6000 km from the main server and 300 km from a copy. Compare the round-trip travel times.
Far server: 12 000 km ÷ 200 000 km/s = 0.06 s = 60 ms. Near copy: 600 km ÷ 200 000 = 0.003 s = 3 ms. The nearby replica is 20 times quicker (travel time only).
Common mistakes
- Thinking "distributed" just means "more computers". They must work together as one service.
- Thinking a load balancer stores the data. It only directs requests.
- Thinking copies are always identical at every moment. After a change, copies may differ for a short time.
- Thinking peer-to-peer has no rules. Peers still follow a shared protocol.