Skip to main content

Track 1: Foundations

Before diving into consensus protocols and complex architectures, you must understand the fundamental challenges that make distributed systems hard. This chapter is the foundation everything else rests on — if your mental model of networks, time, and failure is wrong, every design decision you make downstream will be subtly broken.
Track Duration: 28-38 hours
Modules: 5
Key Topics: Network fundamentals, Time & Ordering, Failure Models, CAP/PACELC

Module 1: Why Distributed Systems?

The Need for Distribution

Every successful system eventually outgrows a single machine:

Types of Distributed Systems

Compute Clusters

Purpose: Process data across many machinesExamples:
  • Hadoop/Spark clusters
  • Kubernetes pods
  • AWS Lambda fleet

Storage Systems

Purpose: Store and retrieve data across machinesExamples:
  • Distributed databases (Cassandra, DynamoDB)
  • Object stores (S3, GCS)
  • File systems (HDFS, Ceph)

Message Systems

Purpose: Enable communication between distributed componentsExamples:
  • Message queues (Kafka, RabbitMQ)
  • Pub/sub systems (SNS, Pub/Sub)
  • Event streams (Kinesis, Event Hubs)

The Eight Fallacies of Distributed Computing

Every distributed systems engineer must internalize these. Peter Deutsch (with additions by James Gosling) identified these assumptions that developers new to distributed systems unconsciously make. Each one is a trap — code that works perfectly on your laptop will fail in production because of one or more of these false beliefs:
Reality: Networks fail constantly.
Defense: Implement retries, timeouts, circuit breakers, and idempotency.
Reality: Every network call has latency.Impact: A service with 10 sequential remote calls adds 500-2000ms latency.Defense: Cache aggressively, parallelize calls, use async processing.
Reality: Bandwidth is expensive and limited.
Defense: Compress data, use efficient serialization (protobuf), batch requests.
Reality: Assume the network is compromised.Attack vectors:
  • Man-in-the-middle attacks
  • DNS spoofing
  • BGP hijacking (internet routing attacks)
  • Insider threats
Defense: mTLS everywhere, zero-trust architecture, encrypt at rest and in transit.
Reality: Network topology changes constantly.
Defense: Use service discovery, health checks, load balancer draining.
Reality: Multiple teams, policies, and even companies.In a typical microservices architecture:
  • Platform team manages Kubernetes
  • Each service team manages their services
  • Security team manages policies
  • Network team manages infrastructure
Defense: Clear ownership, documented interfaces, SLAs between teams.
Reality: Data transfer costs real money.High-traffic example: 1PB egress = $90,000/monthDefense: Keep data close to compute, use CDNs, compress aggressively.
Reality: Different hardware, protocols, and vendors everywhere.
Defense: Abstract hardware differences, use consistent protocols, test on diverse environments.

Module 2: Network Fundamentals

TCP Guarantees and Failures

TCP provides order and reliability, but not much else:

Network Partitions

A network partition occurs when nodes can’t communicate with each other:
Key insight: During a partition, you must choose between:
  • Availability: Both partitions continue serving requests (might diverge)
  • Consistency: Reject requests from the minority partition
Think of it like two bank branches whose phone line gets cut. Do they keep processing transactions independently (available but potentially inconsistent — both might approve overdrafts) or do they stop until the phone line is restored (consistent but unavailable)? There is no option where both branches stay open AND guarantee they agree on the account balance. This is the CAP theorem in a nutshell.

Message Delivery Semantics

Fire and forget - no retries
Use case: Metrics, logs (losing some is acceptable)Risk: Message loss

Failure Detection

How do you know if a node is dead?
┌─────────────────────────────────────────────────────────────────────────────┐ │ PHYSICAL CLOCK PROBLEMS │ ├─────────────────────────────────────────────────────────────────────────────┤ │ │ │ CLOCK SKEW (different machines show different times) │ │ ──────────── │ │ Machine A: 10:00:00.000 │ │ Machine B: 10:00:00.150 (150ms ahead) │ │ Machine C: 09:59:59.850 (150ms behind) │ │ │ │ NTP SYNCHRONIZATION │ │ ─────────────────── │ │ NTP accuracy: 1-10ms on LAN, 10-100ms over internet │ │ During sync: Clock can jump forward or backward! │ │ │ │ CLOCK DRIFT │ │ ──────────── │ │ Quartz crystals drift: ~50ppm (50 microseconds per second) │ │ After 1 hour: 180ms drift possible │ │ │ │ LEAP SECONDS │ │ ──────────── │ │ Added to compensate for Earth’s rotation slowing │ │ 23:59:59 → 23:59:60 → 00:00:00 │ │ Famously caused issues at Reddit, LinkedIn, Mozilla (2012) │ │ │ └─────────────────────────────────────────────────────────────────────────────┘

Vector Clocks

Vector clocks track causality - they tell you if two events are related: Vector Clocks

Advanced: Vector Clocks vs. Version Vectors

In many textbooks, “Vector Clocks” and “Version Vectors” are used interchangeably. However, at the Staff level, you must distinguish between them:
  1. Vector Clocks (Causality): Track the relationship between events in a distributed system. Every interaction (send/receive) increments the clock. They are used to build a partial order of all operations.
  2. Version Vectors (Conflict Tracking): Track the relationship between versions of data. They only increment when data is updated. They are used in systems like DynamoDB and Riak to detect if two versions of a document are in conflict (siblings).
Staff Tip: If you only need to know if two versions of a file conflict, use Version Vectors. If you need to know if a message was “caused” by another message (e.g., in a distributed debugger), use Vector Clocks.

Implementation:

Hybrid Logical Clocks (HLC)

Combines physical and logical clocks:

TrueTime (Google Spanner)

Google’s hardware-based approach:

Module 4: Failure Models

Understanding what can go wrong is crucial for building resilient systems.

Types of Failures

Node crashes and stays downCharacteristics:
  • Clean failure, detectable
  • Node doesn’t recover with corrupted state
  • Easiest to handle
Example: Hardware failure, power loss

Partial Failures

The defining challenge of distributed systems:

Gray Failures

Gray failures are the distributed systems equivalent of a slow gas leak — everything seems fine until suddenly it is not. They are far more common and insidious than clean crash failures, and they are the cause of most prolonged production outages:
Production wisdom: Gray failures are why you must monitor percentile latencies (p99, p99.9), not just averages or error rates. A node with a failing disk might serve 99% of requests normally but add 10 seconds of latency to the rest. The average latency barely moves, but 1% of your users have a terrible experience. Your dashboards should make this visible.

Byzantine Fault Tolerance (BFT)

Byzantine failures are the most challenging to handle:
PBFT Algorithm Overview:
When to Use BFT:
  • Blockchain and cryptocurrency systems
  • Multi-party financial transactions
  • High-security government systems
  • Any system with mutually distrustful participants
Performance Trade-off: BFT protocols are expensive! PBFT requires O(n²) messages per consensus round. Most internal systems use simpler crash-fault tolerant protocols (Raft, Paxos) which only require 2f+1 nodes.

Network Partition Simulation

Understanding how to test partition tolerance:

Chaos Engineering for Failure Testing

Best Practice: Don’t wait for production failures to discover weaknesses. Proactively inject failures in controlled environments.

Module 5: CAP and PACELC Theorems

CAP Theorem

CAP is often misunderstood. It only applies during a network partition.
CAP Theorem

CP vs AP Deep Dive

CP Systems

During partition, choose ConsistencyBehavior:
  • Minority partition stops accepting writes
  • May reject reads too (for linearizability)
  • Majority partition continues
Good for:
  • Financial systems
  • Inventory management
  • Any system where stale data is dangerous
Example: Bank account balance

AP Systems

During partition, choose AvailabilityBehavior:
  • Both partitions continue serving
  • May return stale/conflicting data
  • Resolve conflicts after partition heals
Good for:
  • Social media feeds
  • Shopping carts
  • DNS
Example: Twitter timeline

PACELC: The Better Framework

CAP only describes behavior during partitions. PACELC adds normal operation:

Consistency Spectrum

Consistency is not binary - it’s a spectrum:

External Consistency vs. Linearizability

For Staff/Principal level engineering, you must distinguish between system-local and global-time guarantees.

Linearizability (Atomic Consistency)

Linearizability is a guarantee for a single object (or key). It ensures that:
  1. Every operation appears to take effect atomically at some point between its invocation and response.
  2. All operations are seen in the same order by all participants.
  3. If a write WW completes before a read RR starts (according to the system’s observer), RR must see the result of WW.

External Consistency (Strict Serializability)

External Consistency is a stronger, global guarantee, famously used by Google Spanner. It combines Linearizability with Serializability across multiple shards/objects and respects global wall-clock time. The “Why”: In a multi-shard system, two transactions T1T_1 and T2T_2 might affect different shards. If T1T_1 finishes at 10:00:01 and a client then starts T2T_2 at 10:00:02, T2T_2 must see the effects of T1T_1. Without external consistency (e.g., using only local Raft groups with skewed clocks), T2T_2 might get a timestamp that is logically earlier than T1T_1, leading to causal violations.

Module 6: Modern Infrastructure & Serverless Internals

To achieve the “Principal” level of understanding, one must look beyond the logical protocols and into the physical isolation models that power modern clouds. The shift from monolithic VMs to Micro-VMs has fundamentally changed how we scale distributed systems.

The Serverless Paradox

Serverless (AWS Lambda, Google Cloud Functions) promises “no servers,” but in reality, it requires the most complex server management in existence. The challenge is The Cold Start Problem vs. Isolation Security.
  • Traditional VMs (EC2): Strong isolation, but slow to boot (30s+). Too slow for “on-demand” execution.
  • Containers (Docker/K8s): Fast to start, but weak isolation (shared kernel). Dangerous for multi-tenant code execution.

Micro-VM Internals: AWS Firecracker

AWS Firecracker is the technology behind Lambda and Fargate. It uses KVM (Kernel-based Virtual Machine) but stripped down to the bare essentials.
  1. Minimal Device Model: Firecracker drops legacy hardware support (no USB, no video). It only provides:
    • Net (Network)
    • Block (Storage)
    • Serial (Console)
    • Balloon (Memory management)
  2. Snapshot-and-Restore: Instead of booting a kernel, Firecracker can restore from a Memory Snapshot. This reduces “boot” time from seconds to ~10ms.
  3. Jailer: A secondary security layer that uses cgroups, namespaces, and seccomp to ensure that even if a Micro-VM is compromised, it cannot touch the host or other VMs.

Implications for Distributed Design

When designing systems for serverless:
  • Statelessness is Mandatory: Because the Micro-VM can be killed or recycled at any moment.
  • Sidecar Overhead: Traditional sidecars (Service Mesh) add too much latency for 50ms functions. We move towards Proxyless Mesh or Library-based approaches.
  • Connection Pooling: Traditional database connection pools fail at serverless scale. We use RDS Proxy or HTTP-based database protocols (Data API).

Key Interview Questions

Answer: No, provably impossible (FLP theorem). But we can achieve it practically by:
  • Using timeouts (introduces synchrony assumption)
  • Randomization
  • Failure detectors
Most practical systems (Raft, Paxos) assume partial synchrony.
Answer: You can’t with certainty. Approaches:
  • Timeouts: Simple but prone to false positives
  • Heartbeats: Regular health checks
  • Phi accrual detector: Probabilistic suspicion level
  • Lease-based: Node must renew lease to be considered alive
In practice, use adaptive timeouts based on historical latency.
Answer:
  • Lamport: Single counter, preserves “happens-before” one direction only
  • Vector clocks: One counter per node, can detect concurrent events
Use Lamport when you just need ordering. Use vector clocks when you need to detect conflicts.
Answer: Frame it as a trade-off discussion:
  • First, acknowledge CAP only applies during partitions
  • Discuss what consistency level you actually need
  • Explain PACELC for normal operation
  • Give examples: “For our payment system, we’re CP because incorrect balance is unacceptable. For user preferences, we’re AP because eventual consistency is fine.”

Next Steps

Continue to Track 2: Consensus Protocols

Learn Paxos, Raft, and other consensus algorithms that power distributed systems

Interview Deep-Dive

Strong Answer:
  • FLP proves that no deterministic algorithm can guarantee consensus in a fully asynchronous system where even one process may crash. The key word is “guarantee” — it is an impossibility of absolute liveness, not of safety.
  • Practical systems like Raft and Paxos work around FLP by introducing partial synchrony assumptions. Specifically, they use timeouts to detect suspected failures. Once you add a timeout, you are no longer in a purely asynchronous model — you are assuming that messages will eventually be delivered within some bound. FLP does not apply to partially synchronous systems.
  • The trade-off is that during periods of true asynchrony (extreme network delays, GC pauses that exceed the timeout), these protocols may fail to make progress (liveness violation) but they never violate safety (they never agree on two different values). This is exactly the right trade-off for production systems: it is acceptable to be temporarily unavailable, but it is never acceptable to corrupt data.
  • Another approach is randomized consensus (like Ben-Or’s protocol), which circumvents FLP by using randomization rather than determinism. The expected number of rounds is finite, but no single execution is guaranteed to terminate.
Follow-up: Can you describe a scenario where the FLP result actually manifests in a production system?Consider a Raft cluster where the leader experiences a long GC pause. During the pause, followers time out and start an election. But the GC pause ends and the old leader resumes sending heartbeats before the election completes. Followers receive heartbeats from the old leader and reset their election timers, canceling the election. If this cycle repeats — GC pause, election starts, GC ends, election cancels — the cluster can enter a livelock where no leader is elected and no progress is made. This is FLP manifesting in practice. The mitigation is randomized election timeouts in Raft (150-300ms range), which break symmetry and make it statistically improbable for this livelock to persist. But “statistically improbable” is not “impossible” — FLP reminds us that no amount of engineering eliminates this possibility entirely.
Strong Answer:
  • A full partition splits the cluster into two groups that cannot communicate at all. This is the classic CAP scenario: each side must independently decide whether to keep serving (AP) or stop (CP). Raft handles this cleanly — the majority side elects a leader, the minority side stops accepting writes.
  • A partial partition is when some nodes can communicate with some but not all other nodes. For example, in a 5-node cluster, node A can talk to B and C, node D can talk to C and E, but A cannot talk to D or E. Node C acts as a bridge. This is more insidious because quorum-based systems may still form quorums that do not overlap correctly, and different nodes have different views of who is reachable.
  • An asymmetric partition is when A can send to B but B cannot send to A. TCP will eventually detect this (via ACK timeouts), but UDP-based gossip or heartbeat systems may not. Node A thinks B is alive (because A’s sends succeed), while B thinks A is dead (because B never receives anything).
  • The distinction matters because most distributed systems are designed for full partitions, but partial and asymmetric partitions are more common in practice (misconfigured firewalls, flaky switches, half-duplex network failures). If your failure detector assumes symmetric communication, an asymmetric partition can cause split-brain.
Follow-up: How would you test your system’s behavior under a partial partition?I would use a network simulation layer (like Jepsen’s iptables-based nemesis or Toxiproxy) that can selectively drop traffic between specific pairs of nodes. The test scenario would create a partition where nodes A-B-C can communicate, D-E can communicate, but C is the only bridge between the two groups. Then I would verify: can the system still make progress? Does it correctly detect the partial partition? Does it degrade gracefully? I would also test the asymmetric case by dropping packets in one direction only using iptables INPUT/OUTPUT rules. The key metrics to watch are: leader election stability, write latency, and most importantly, whether any safety invariants (like linearizability) are violated during the partition.
Strong Answer:
  • I would respectfully push back. BFT is the right tool for a very specific threat model: mutually distrustful participants who may actively lie (blockchains, multi-party financial systems). For the vast majority of internal distributed systems, crash-fault tolerance (CFT) is sufficient and dramatically cheaper.
  • The cost difference is significant. CFT protocols like Raft need 2f+1 nodes to tolerate f failures, with O(n) message complexity per consensus round. BFT protocols like PBFT need 3f+1 nodes with O(n-squared) message complexity. For a 5-node Raft cluster tolerating 2 failures, you need 5 nodes. For BFT tolerating 2 failures, you need 7 nodes, each exchanging quadratically more messages.
  • In practice, the failures that cause production outages are crashes, network partitions, disk failures, and software bugs — not malicious nodes. If you are worried about buggy nodes sending corrupt data, you can add checksums and input validation at the application layer for a fraction of the cost of full BFT.
  • The exception: if you are building a system where different organizations contribute nodes and have economic incentives to cheat, BFT is necessary. This is why blockchain consensus uses BFT variants.
Follow-up: What about the case where a software bug causes a node to send corrupt or inconsistent messages? Is that not effectively Byzantine behavior?This is an excellent point and a real concern. A node with a memory corruption bug or a deserialization error might send valid-looking but incorrect messages — this is sometimes called an “accidental Byzantine” fault. However, the practical response is not to deploy full BFT. Instead, I would use defense-in-depth: checksums on all messages to detect corruption, assertion-based crash detection (if a node detects its own state is inconsistent, it crashes rather than propagating corruption), and Merkle-tree-based anti-entropy to detect data divergence between replicas. These approaches catch accidental Byzantine behavior at a fraction of the cost. For true Byzantine resilience, I would only invest in it for the specific data paths where corruption would be undetectable and catastrophic.
Strong Answer:
  • PACELC extends CAP by asking: even when there is no partition, what trade-off does the system make between latency and consistency? The framework is: if Partition, choose Availability or Consistency; Else, choose Latency or Consistency.
  • PA/EL (Available during partitions, Low latency normally): DynamoDB and Cassandra. During a partition, they keep serving from whichever replicas are reachable. During normal operation, they return responses from the nearest replica without waiting for quorum, giving low latency but eventual consistency.
  • PA/EC (Available during partitions, Consistent normally): MongoDB with default read concern. During normal operation, reads go to the primary, which is consistent. During a partition, secondaries can still serve reads (with possible staleness), maintaining availability.
  • PC/EL (Consistent during partitions, Low latency normally): Yahoo’s PNUTS (now largely historical). It blocks writes during partitions to maintain per-record consistency, but during normal operation it routes reads to the nearest replica for low latency using “timeline consistency.”
  • PC/EC (Consistent during partitions, Consistent normally): Google Spanner. It blocks during partitions (choosing consistency), and during normal operation it uses TrueTime and commit-wait to provide external consistency, accepting higher latency (1-7ms commit wait) for global ordering guarantees.
Follow-up: Can a single system offer different PACELC trade-offs for different operations?Absolutely, and this is how most sophisticated production systems work. Cassandra is the canonical example: you can set consistency level per query. A write with QUORUM and a read with QUORUM gives you PC/EC behavior for that specific operation. A write with ONE and a read with ONE gives you PA/EL behavior. The same cluster, the same data, but different consistency-latency trade-offs per request. Spanner does something similar with its read-only transactions: a “strong read” gives you external consistency (PC/EC), while a “stale read” with a bounded staleness of 10 seconds gives you lower latency (closer to PC/EL). The design insight is that consistency should be a dial, not a switch — and different use cases within the same application should be able to turn that dial independently.