Introduction to Distributed Systems

What a distributed system is, why we need one, and the properties to reason about.

Distributed Systems Essentials

Introduction to Distributed Systems

A distributed system is a collection of independent but coherent nodes (computing elements like machines or processors). These nodes interact with each other to serve the functionalities of an application. To the user, however, it seamlessly appears as a single unified system.

Why Do We Need Distributed Systems?

  • Avoid Single Point of Failure: If an application relies on a single node and that node fails, the entire system goes down. Distributed systems use replicas and backups to keep the application running during node failures.
  • Handle Large Traffic: A single machine has strict hardware limitations (RAM, CPU, storage). When scaling to thousands or millions of concurrent users, a single system cannot handle the load.
  • Reduce Latency: Distributing traffic and data prevents performance bottlenecks, ensuring users get fast responses.

Key Properties of Distributed Systems

1. Availability Availability refers to the "uptime" of a system—the percentage of time the system is operational compared to the total time.

  • 99% Availability: System is down for about 4 days a year.
  • 99.999% (Five Nines) Availability: System is down for only ~5.26 minutes a year (often targeted by companies like Amazon).

Important Availability Terms:

  • SLA (Service Level Agreement): The legal promise or contract made by the system providers to the clients.
  • SLO (Service Level Objective): The individual, specific targets within the SLA (e.g., "Latency will be under 20ms").
  • SLI (Service Level Indicator): The actual, measured performance of the system in the real world.

2. Reliability Reliability means the system can continue to operate correctly even during partial failures.

  • Hardware Faults: Disk failures or power outages. These are handled by data replication and backup power grids.
  • Software Faults: Bugs, bad code, or unhandled errors. These are handled by unit testing, integration cycles, and strict error handling (try/except blocks).

3. Scalability Scalability is the ability to handle increased load seamlessly.

  • Vertical Scaling: Upgrading the resources (more RAM, bigger hard drive) on a single machine. This gets very expensive and has strict physical limits.
  • Horizontal Scaling: Adding more machines (nodes) to the system to distribute the load. This is the industry standard.

Load Balancing

Load balancing is the process of distributing incoming traffic evenly across multiple nodes so that no single machine is overwhelmed.

Understanding System Load Before balancing, the system needs to measure the load using metrics like:

  • QPS (Queries Per Second): The volume of requests hitting the servers, which fluctuates throughout the day.
  • Read-to-Write Ratio: Figuring out if the system is "Read-heavy" (like a blog) or "Write-heavy" (like YouTube video uploads).
  • Percentiles of Response Time: Analyzed using Pareto charts. For example, if the 99th percentile latency is 200ms, it means 99% of user requests are served in under 200ms.

The Role of a Load Balancer

  • Evenly divides requests among available nodes.
  • Performs constant "status checks" to ensure nodes are healthy.
  • Stops sending traffic to dead or failing nodes.
  • Triggers the initiation of new nodes if the current ones are overwhelmed. (Note: Load balancers themselves must be replicated and mapped to DNS to avoid becoming a single point of failure).

Load Balancing Algorithms

Load balancers use specific algorithms to decide where to send incoming traffic:

  • Hashing: Converts a Session ID or User ID into a number, takes the modulo of the number of nodes, and routes the traffic. Ensures the same user hits the same server.
  • Endpoint Routing: Directs traffic based on the URL (e.g., /courses traffic goes to Node A, /newsletter traffic goes to Node B).
  • Round Robin: Distributes requests sequentially in a continuous cycle (Node 1, then Node 2, then Node 3, repeat).
  • Least Pending Requests: Routes traffic to the node that currently has the fewest active requests waiting in its queue.
  • Least Loaded: Routes traffic to the node currently utilizing the least amount of CPU/RAM.
  • Random Routing: Routes traffic randomly based on calculated probability.