📘 CodingMarble Learn

Distributed Systems: Sharing Work and Data in a Network

A distributed system is many computers in a network working together as one service. They share the work (many servers, a load balancer) and the data (copies or pieces on different machines). This makes a service faster and able to survive a failure, but needs care to keep the copies in step.

🎬 Step-by-step story

  1. Phones (clients) ask one big computer, the server, for pages and data. The server answers each request. Few users: no waiting.
  2. Too many users! One server can answer only a few requests each second, so the rest wait in a line (the yellow blocks).
  3. Share the work: three servers and a traffic manager (load balancer). It sends each request to the least busy server, so every line stays short.
  4. Each server keeps a copy of the data (blue disc). One server breaks, but the others still answer. This is replication.
  5. Another way: no boss computer. Peers swap pieces of a file with each other until all have everything.
  6. Your turn: change the requests and the servers, and break one server. Watch the waiting line and the waiting time.

Tip: drag the 3D scene to turn it. Use two fingers to zoom.

🤔 Common doubts, cleared

What exactly is a server?

A server is just a computer (or program) that waits for requests and answers them. The big blue boxes in step 1 are servers.

Why does a line form at the server?

Requests arrive faster than the server can answer them (step 2 shows 5 per second against about 1 answered), so the extra ones wait. Watch the yellow blocks pile up.

Does the load balancer do the work itself?

No. It only chooses which server gets each request. Servers still do the real work, as you see in step 3.

If a server breaks, is the data lost?

Not if copies exist on other servers (the blue discs). That is why we use replication. Check the broken grey server in step 4.

How can a network work without a central server?

Each peer holds some pieces and shares them. Together they hold the whole file. Step 5 shows pieces moving until all peers are complete.

Is more servers always better?

No. More servers cost more and make copies harder to keep in step. Use the sliders in step 6 to find how many are really needed for a given load.

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:

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

Key formulas and definitions

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

Practice quiz

1. In client-server, who asks for the data?
2. What does a load balancer do?
3. Keeping copies of the same data on several machines is called:
4. A single point of failure means:
5. In peer-to-peer, each computer is:

Practice: answer these yourself

Type or choose your answer, then press Check. Use a hint if you are stuck; the full solution appears after you answer.

Frequently asked questions

What is a distributed system in simple words?

Many computers connected by a network that share the work and data and act like one service for the user.

What is the difference between client-server and peer-to-peer?

In client-server some computers (servers) only serve and others (clients) only ask. In peer-to-peer every computer does both and there is no central boss.

Why do big websites use many servers?

To answer more users at the same time, to keep working when a machine breaks, and to keep data close to users.

Where this is taught

NetherlandsHAVO 5 (eindexamenjaar)Elective theme: Networks
NetherlandsVWO 6 (eindexamenjaar)Elective theme: Networks

Learn first

Learn next

Related lessons

All Computer Science lessons