Skip to content

Distributed Systems Fundamentals

Core distributed systems concepts that appear across all top-tier engineering interviews.

You can only guarantee two of three:

  • Consistency: Every read sees the most recent write
  • Availability: Every request gets a response
  • Partition Tolerance: System works despite network partitions

In practice, partitions happen, so you’re choosing between CP (consistent but may reject requests) and AP (available but may return stale data).

SystemTypeTrade-off
ZooKeeperCPRejects writes during partition
CassandraAPMay return stale data, eventual consistency
SpannerCPUses TrueTime for global consistency, may increase latency
DynamoDBTunableChoose between strong and eventual consistency per read

From strongest to weakest:

ModelGuaranteePerformance
LinearizabilityEvery read sees the latest write globallySlowest
Sequential consistencyAll processes see operations in same orderSlow
Causal consistencyCausally related operations are orderedMedium
Eventual consistencyAll replicas converge eventuallyFastest
  • Linearizability: Bank account balance, leader election, distributed locks
  • Causal consistency: Social media feeds, collaborative editing
  • Eventual consistency: DNS, caches, analytics counters, recommendations
[Leader] ---AppendEntries--> [Follower 1]
| [Follower 2]
| [Follower 3]
| [Follower 4]
|
+--- Committed when majority (3/5) acknowledges

Leader election:

  1. Follower times out (no heartbeat)
  2. Becomes candidate, votes for self, requests votes
  3. Wins with majority, becomes leader
  4. Sends heartbeats to maintain authority

Log replication:

  1. Leader appends entry to log
  2. Sends to all followers
  3. Committed when majority acknowledges
  4. Applied to state machine

Key properties:

  • At most one leader per term
  • Committed entries survive leader changes
  • Simple enough to implement correctly
  • More general than Raft, harder to understand
  • Proposer, Acceptor, Learner roles
  • Two-phase: Prepare -> Accept
  • Multi-Paxos for repeated consensus (similar to Raft)
[Client] --> [Leader] --> [Follower 1]
--> [Follower 2]
--> [Follower 3]
  • All writes go through leader
  • Followers replicate asynchronously (or synchronously for strong consistency)
  • Failover: Promote a follower to leader if leader fails
[Client A] --> [Leader DC1] <--> [Leader DC2] <-- [Client B]
  • Multiple leaders accept writes (e.g., one per data center)
  • Conflict resolution needed: last-writer-wins, merge, custom resolution
  • Use case: Multi-region deployment for low write latency
[Client] --> [Node 1] (W)
--> [Node 2] (W)
--> [Node 3] (W)
  • Client sends writes to multiple nodes
  • Quorum: W + R > N ensures read sees latest write
    • N=3, W=2, R=2: Strong consistency
    • N=3, W=1, R=1: Fastest, eventual consistency
  • Anti-entropy: Background process syncs replicas
StrategyDescriptionProsCons
Hash partitioninghash(key) % NEven distributionRange queries span all partitions
Range partitioningKey ranges per partitionEfficient range queriesHot spots possible
Directory-basedLookup table maps key to partitionFlexibleLookup table is bottleneck
Consistent hashingHash ring with virtual nodesEasy rebalancingLess control over placement
  • Add random suffix to hot keys (scatter reads/writes)
  • Replicate hot partitions
  • Application-level caching for hot data
  • Dynamic partition splitting
Phase 1 (Prepare):
Coordinator --> "Can you commit?" --> Participant A: "Yes"
--> Participant B: "Yes"
Phase 2 (Commit):
Coordinator --> "Commit" --> Participant A: Done
--> Participant B: Done

Problem: Blocking. If coordinator fails after sending “prepare” but before “commit”, participants are stuck holding locks.

For long-running distributed transactions:

Step 1: Create Order | Compensate: Cancel Order
Step 2: Reserve Payment | Compensate: Release Payment
Step 3: Reserve Inventory| Compensate: Release Inventory
Step 4: Ship | Compensate: Return

If any step fails, execute compensating actions in reverse. Non-blocking, eventually consistent.

  • Choreography: Each service emits events, next service reacts
  • Orchestration: Central coordinator drives the saga steps
AlgorithmDescriptionUse Case
Round RobinDistribute sequentiallyEqual-capacity servers
Weighted Round RobinProportional to capacityMixed server sizes
Least ConnectionsSend to least-busy serverVariable request duration
Consistent HashingHash-based routingStateful services, caching
RandomRandom selectionSimple, surprisingly effective
  • Layer 4 (TCP): Fast, low overhead, no content inspection
  • Layer 7 (HTTP): Can route based on URL, headers, cookies. More flexible.
StrategyRead PathWrite PathConsistency
Cache-asideApp checks cache, misses hit DBApp updates DB, invalidates cacheEventual
Read-throughCache fetches from DB on missApp writes to cache, cache writes to DBEventual
Write-throughRead from cacheWrite to cache and DB simultaneouslyStrong
Write-behindRead from cacheWrite to cache, async write to DBWeak

The hardest problem in computer science (after naming things):

  • TTL: Simple, but data can be stale up to TTL
  • Event-based: Invalidate on write events. Low staleness but complex.
  • Version-based: Include version in cache key. Never stale but requires version tracking.
[Client] --> [Browser Cache] --> [CDN] --> [API Gateway Cache]
--> [Application Cache (Redis)]
--> [Database]
PatternDescriptionDeliveryUse Case
Point-to-PointOne producer, one consumerAt-most-once or at-least-onceTask processing
Pub/SubOne producer, many consumersFan-outEvent notification
Competing ConsumersMany consumers, each message processed onceAt-least-onceLoad distribution
Event SourcingLog of all state changesExactly-once (with dedup)Audit trail, replay
  • At-most-once: Send and forget. May lose messages. Fastest.
  • At-least-once: Retry until acknowledged. May duplicate. Most common.
  • Exactly-once: Deduplication + idempotent processing. Hardest, slowest.

In practice, design for at-least-once delivery + idempotent consumers.

AlgorithmDescriptionProsCons
Token BucketTokens added at rate R, consumed per requestAllows bursts, smoothState per client
Leaky BucketRequests drain at fixed rateSmooth outputNo burst tolerance
Fixed WindowCount requests per time windowSimpleBoundary burst (2x at window edge)
Sliding Window LogTrack timestamp of each requestAccurateMemory-intensive
Sliding Window CounterWeighted average of current and previous windowLow memory, accurateApproximate
Option 1: Centralized (Redis)
[Client] --> [API Server] --> [Redis: INCR + EXPIRE] --> Allow/Deny
Option 2: Local + Sync
Each server has local counter, periodically syncs to central store
Allows slightly over limit but no network hop per request
PatternLatencyCouplingReliability
Synchronous RESTLowHighRequest fails if service down
Synchronous gRPCLowerHighSame, but with protobuf efficiency
Async messagingHigherLowQueue buffers if service down
Event-drivenVariableLowestMost resilient
  • Circuit Breaker: Stop calling failing services, fail fast
  • Retry with Backoff: Exponential backoff + jitter
  • Bulkhead: Isolate resources per dependency (separate thread pools)
  • Timeout: Always set timeouts. Never wait indefinitely.
  • Fallback: Degrade gracefully (cached response, default value)
[Service A] --> [Sidecar Proxy A] --> [Sidecar Proxy B] --> [Service B]
| |
[Control Plane (Istio/Linkerd)]

Handles: Load balancing, mTLS, retries, circuit breaking, observability — without application code changes.

PillarWhat It CapturesTools
LogsDiscrete eventsELK, Loki, CloudWatch
MetricsAggregated measurementsPrometheus, Datadog, CloudWatch
TracesRequest flow across servicesJaeger, Zipkin, X-Ray
  • Rate: Requests per second
  • Errors: Error rate (percentage of failed requests)
  • Duration: Latency distribution (p50, p95, p99)
SLI (Indicator): p99 latency < 200ms
SLO (Objective): 99.9% of requests meet the SLI over 30 days
Error Budget: 0.1% of requests can violate (43.2 minutes/month)

When error budget is exhausted: Freeze deployments, focus on reliability.

OperationTime
L1 cache reference1 ns
L2 cache reference4 ns
Main memory reference100 ns
SSD random read16 μs
HDD seek4 ms
Send 1 KB over 1 Gbps network10 μs
Read 1 MB from memory250 μs
Read 1 MB from SSD1 ms
Read 1 MB from HDD20 ms
Roundtrip same datacenter0.5 ms
Roundtrip CA -> Netherlands150 ms
1 day = 86,400 seconds ≈ 100K seconds
1 million requests/day ≈ 12 requests/second
1 billion requests/day ≈ 12,000 requests/second
1 KB * 1 billion = 1 TB
1 MB * 1 million = 1 TB