Week 04 · Research Update

When One Computer Is Not Enough: The Fascinating World of Distributed Systems

1. What happens when the problem you want to solve becomes too big for any one computer?

Three distinct walls tend to appear:

  • The compute wall: A single machine has a hard ceiling on how many operations it can perform per second. Training a large machine learning model requires more raw arithmetic than any single processor can deliver within a useful timeframe. The only way past this ceiling is to divide the arithmetic across many processors working in parallel.
  • The memory/storage wall: Some datasets are simply too large to fit on one machine's disk, let alone in its RAM. A web-scale search index, a global genomic database, or the full transaction history of a large bank cannot be held by a single server. The data itself has to be partitioned across machines.
  • The reliability wall: Even when a single machine could do the job, doing so is risky: if a task takes three weeks to run and the machine fails on day twenty, all that work is lost. Distributing the task across multiple machines means a single hardware failure costs you a fraction of the work, not all of it.

Once any of these three walls is hit, the problem becomes a genuinely distributed-systems problem: now you must design for machines that communicate over unreliable networks, that can fail independently, and that need to agree on a shared view of the work.

2. Suppose 1,000 computers work together. Do you now have one computer that is 1,000 times more powerful? Why or why not?

No, and this is one of the foundational lessons of Amdahl's Law (Amdahl, 1967).

  • Sequential bottlenecks limit speedup: Any real task has some portion that must happen in strict order (setup, coordination, final aggregation of results) and some portion that can be split across machines. Amdahl's Law shows that if even 10% of a task is inherently sequential, the maximum possible speedup, no matter how many machines you throw at it, caps out at 10x, nowhere near 1,000x. The sequential fraction becomes the dominant cost as parallelism increases.
  • Coordination has a real cost: Machines working together must communicate: sending messages, waiting for acknowledgments, and synchronizing shared state. This communication consumes time, bandwidth, and CPU cycles that a single machine never had to spend on itself.
  • New failure modes appear: A single machine either works or it doesn't. A thousand machines introduce network latency, partial failures (some machines are up, some are down, some are slow), and clock disagreement between nodes; none of which exist in a single-machine world.

So while adding machines genuinely helps up to a point, returns diminish and can even go negative if coordination overhead grows faster than the useful work being parallelized. This is precisely why just adding more servers is never a free lunch, and why distributed systems engineering is a discipline in its own right rather than simple hardware shopping.

3. Can 1,000 computers agree on something if some of them fail or even lie?

This is the founding question of an entire subfield: distributed consensus.

When machines can only fail silently (crash and stop responding, but never send false information), a family of consensus algorithms solves the problem. The two most famous are:

  • Paxos: introduced by Leslie Lamport (1998), the original and mathematically rigorous (but notoriously difficult to fully understand) protocol for getting a group of machines to agree on a single value even as some of them crash.
  • Raft: designed by Diego Ongaro and John Ousterhout specifically as a more understandable alternative to Paxos. Raft is a consensus algorithm developed at Stanford University, designed to be easier to understand than previous algorithms like Paxos while still providing strong fault tolerance and leader-election capabilities. It works by electing a leader among the servers, who is responsible for managing replication and ensuring all followers hold the same data; if the leader fails, the cluster elects a new one.

When machines can lie (sending different, contradictory, or malicious information to different peers, whether due to a bug or an actual attacker) the problem becomes the much harder Byzantine fault tolerance problem, formalized by Lamport, Shostak, and Pease (1982) in their paper "The Byzantine Generals Problem." Their landmark result: to reach agreement while tolerating up to f lying or arbitrarily faulty machines out of n total, you need at least n ≥ 3f + 1 machines. Below that threshold, no protocol can guarantee correct agreement.

This is exactly the problem that blockchain systems like Bitcoin and Ethereum had to solve at massive, decentralized scale.

A related and equally famous impossibility result constrains what any of these systems can promise: the CAP theorem, first articulated by Eric Brewer (2000) and formally proven by Gilbert and Lynch (2002). It states that a distributed data store can provide only two of three guarantees simultaneously: Consistency (every read gets the most recent write, or an error), Availability (every request gets a non-error response, though not necessarily the latest data), and Partition tolerance (the system keeps working despite dropped or delayed messages between nodes). Since network partitions cannot be avoided in a real distributed system, in practice a system designer must choose between consistency and availability whenever a partition occurs. This is why some databases (traditional relational databases) favor strict consistency, while others (many NoSQL systems) favor staying available even at the cost of temporarily inconsistent answers.

4. When you use ChatGPT, Google, Instagram, or an online game, where is the computation actually happening?

It's spread across a layered hierarchy:

  • My own device: rendering the interface, capturing your keystrokes, decoding video for playback.
  • Edge and CDN (content delivery network) servers: positioned geographically close to me, cache static content so requests don't have to travel across an ocean every time.
  • Regional and global data centers: do the heavy lifting. Google's search index is sharded and replicated across thousands of machines worldwide; Instagram's photo storage, social graph, and recommendation ranking run on distributed databases and machine-learning clusters; a single ChatGPT-style response is produced by a cluster of GPUs coordinating to serve one request, often orchestrated across multiple data centers for redundancy.
  • On-device computation: is increasingly used too. On-device machine learning models handle tasks like face detection in your camera app locally, trading some accuracy for lower latency and better privacy, which is itself a reminder that the boundary of distribution keeps shifting as hardware improves.

The honest, slightly unsettling answer: the computation behind almost any app you use daily is happening in a global network of machines you will never see directly, often in several physical locations simultaneously and dynamically, depending on load, your location, and failure conditions at that exact moment.

5. If you could make millions of computers behave like one dependable machine, what could humanity build that we cannot build today?

Once a distributed system can be trusted the way you trust a single reliable machine, new categories of things become possible:

  • Planet-scale real-time coordination: global financial settlement systems, live worldwide epidemiological modeling that updates as new data arrives, or coordinated disaster response systems that can reroute resources in real time across continents.
  • Radical fault tolerance: services that stay online even if an entire data center is destroyed by a natural disaster.
  • Trustless coordination among strangers: currencies, contracts, or shared records that function correctly without requiring a central authority anyone has to trust.

References

  1. Amdahl, G. M. (1967). Validity of the single processor approach to achieving large scale computing capabilities. AFIPS Conference Proceedings, 30, 483–485. Link
  2. Lamport, L., Shostak, R., & Pease, M. (1982). The Byzantine Generals Problem. ACM Transactions on Programming Languages and Systems, 4(3), 382–401. Link
  3. Ongaro, D., & Ousterhout, J. (2014). In Search of an Understandable Consensus Algorithm (Extended Version). Proceedings of the 2014 USENIX Annual Technical Conference (USENIX ATC '14), Stanford University. Link
  4. Gilbert, S., & Lynch, N. (2002). Brewer's Conjecture and the Feasibility of Consistent, Available, Partition-Tolerant Web Services. ACM SIGACT News, 33(2), 51–59. Link

Further reading and tools

  1. The Raft Consensus Algorithm
  2. CAP theorem, Wikipedia
  3. MIT 6.824 / 6.5840: Distributed Systems
  4. Baeldung: "How the Raft Consensus Algorithm Works"