Learning track 07
Distributed Systems — From First Principles to Production
Twelve modules from partial failure to architect-level playbooks. One ShardPay payment ledger grows with every chapter — every failure mode you'll face at scale.
- Articles
- 54
- Modules
- 12
- Total read
- 9 hr 22 min
- Streak
- 0 days
Start here
What Makes a System Distributed
Distributed Foundations
What makes a system distributed and why partial failure is the default
- What Makes a System DistributedIndependent nodes, partial failure, and no shared clock — the three facts every distributed design must accept.
- The 8 Fallacies of Distributed ComputingWhy assuming reliable networks and zero latency leads to production outages — and how ShardPay designs around each fallacy.
- When to Distribute (and When Not)The operational tax of distribution, when a modulith wins, and how ShardPay knew it was time to split.
- Fault, Failure, and ErrorDetect faults, tolerate failures, and prevent errors from becoming customer-visible incidents.
- Reliability, Availability, and DurabilityNines math, RPO/RTO, and how ShardPay sets SLA targets for payment ledger durability.
Performance and Scale
Latency, throughput, scaling models, and capacity planning
- Latency, Throughput, and Little's LawWhy optimizing average latency hides tail problems — and how L = λW governs ShardPay queue depth.
- Vertical vs Horizontal ScalingScale up until you cannot, then scale out — stateless tiers first, stateful tiers with a plan.
- Stateless vs Stateful TiersWhich ShardPay services can scale freely and which need sticky routing, replication, or partitioning.
- Capacity Planning and HeadroomNever run production at 100% — spikes, failover capacity, and ShardPay's Black Friday playbook.
Networking in Distributed Systems
RTT, DNS, load balancing, CDN, and service discovery
- RTT and the Bandwidth-Delay ProductWhy cross-region transfers are slow even on fast links — and how ShardPay batches ledger sync.
- DNS, Service Discovery, and Health-Aware RoutingHow ShardPay instances find each other and how health checks prevent routing to dead nodes.
- Load Balancing Algorithms and Sticky SessionsRound-robin, least connections, consistent hashing at the edge, and when stickiness hurts scale.
- CDN, Edge Caching, and AnycastPush static assets and read-heavy API responses closer to users — ShardPay's receipt PDF delivery.
Consistency and CAP
CAP, PACELC, consistency models, and quorum design
- CAP Theorem — Theory and PracticeDuring a partition you choose consistency or availability — ShardPay's payment path picks CP.
- PACELC — Beyond CAPEven without partitions, latency forces consistency vs response-time trade-offs.
- Consistency Models (Strong to Eventual)Read-your-writes, monotonic reads, and when ShardPay's balance API can return stale data.
- Linearizability vs SerializabilitySingle-operation vs transaction ordering guarantees — critical for ShardPay double-spend prevention.
- Quorums, Read/Write Concerns, and SLAsW + R > N for strong reads, tunable concerns, and mapping quorum math to customer SLAs.
Replication and Partitioning
Topologies, lag, sharding, consistent hashing, and hot spots
- Replication TopologiesLeader-follower, multi-leader, and leaderless replication — when each fits ShardPay's ledger.
- Read Replicas and Replication LagServe reads from followers but know your staleness budget — ShardPay balance queries.
- Partitioning and Data LocalityRange, hash, and composite keys — how ShardPay shards accounts without cross-shard transactions.
- Consistent Hashing with Virtual NodesAdd or remove nodes without remapping the entire key space — ShardPay's shard ring.
- Hot Spots, Skew, and RebalancingCelebrity accounts and power-law traffic — detect skew and rebalance without downtime.
Coordination and Consensus
Leader election, Raft, locks, gossip, and membership
- Leader Election and FailoverOne writer at a time, clean failover, and fencing stale leaders — ShardPay's shard primaries.
- Raft Consensus — How It WorksLeader election, log replication, and safety — the algorithm behind etcd and many control planes.
- Paxos and ZAB — Conceptual OverviewThe classics behind modern consensus — enough depth to read ZooKeeper and Kafka KRaft docs.
- Distributed Locks and Fencing TokensWhy naive Redis locks fail and how monotonic fencing tokens prevent stale leader writes.
- Gossip Protocols and Cluster MembershipDecentralized health spread without a single coordinator — Cassandra and ShardPay discovery.
Time, Ordering, and Identity
Clocks, logical time, sequencers, and causal consistency
- Physical vs Logical ClocksWall clocks lie across datacenters — why ShardPay never orders events by `System.currentTimeMillis()`.
- Lamport Timestamps and Vector ClocksEstablish happened-before without synchronized clocks — detecting concurrent ShardPay transfers.
- Distributed ID and Sequencer DesignSnowflake IDs, lease-based sequencers, and why ShardPay transfer IDs are globally unique.
- Happens-Before and Causal ConsistencyPreserve cause-effect ordering without paying for linearizability on every read.
Reliability Patterns
Timeouts, retries, circuit breakers, idempotency, and chaos testing
- Timeout Hierarchies and Latency BudgetsPer-hop timeouts must sum under the end-to-end SLA — ShardPay's 200ms transfer budget.
- Retries, Backoff, and JitterRetry storms take down healthy services — exponential backoff with full jitter.
- Circuit Breaker and BulkheadFail fast and isolate thread pools so one slow dependency cannot exhaust ShardPay workers.
- Idempotency Keys and Safe RetriesSame transfer request twice must not debit twice — idempotency keys at the API gateway.
- Chaos Engineering and Jepsen TestingProve your consistency claims under partition — not just in design docs.
Messaging Contracts
Delivery semantics, ordering, outbox, and event coupling
- At-Most, At-Least, and Exactly-Once DeliveryWhat brokers promise vs what your handlers must enforce — ShardPay settlement events.
- Ordering Guarantees Across ServicesPer-partition ordering, total order costs, and when ShardPay accepts out-of-order settlement.
- Outbox and Inbox PatternsAtomically write business state and an event — the bridge between DB commits and Kafka.
- Temporal and Spatial CouplingWhen events decouple teams and when hidden coupling via schemas creates distributed monoliths.
Distributed Transactions
2PC limits, sagas, choreography, and compensation
- Two-Phase Commit and Its LimitsAtomic commit across nodes — why ShardPay avoids 2PC for cross-shard transfers.
- Saga Pattern — OrchestrationA central coordinator drives compensating steps when ShardPay cross-shard transfers fail mid-flight.
- Saga Pattern — ChoreographyServices react to events without a central orchestrator — trade-offs for ShardPay refunds.
- Compensation and Eventual ConsistencyDesign reversible operations and reconcile ledgers when perfect atomicity is impossible.
Observability and Operations
Tracing, SLIs, golden signals, and debugging at scale
- Distributed Tracing — Context PropagationTrace IDs across ShardPay services — follow one transfer through gateway, ledger, and fraud.
- SLI, SLO, and Error BudgetsMeasure what users feel, set targets, and spend error budget on risky launches.
- Golden Signals and Tail LatencyLatency, traffic, errors, saturation — why p99 matters more than average for ShardPay.
- Debugging Distributed FailuresCorrelate traces, logs, and metrics to find the one slow hop in a 12-service transfer.
Expert Architecture Playbook
Failure domains, cells, multi-region, team topology, and interviews
- Failure Domains and Blast RadiusIsolate AZs, regions, and cells so one incident cannot drain every ShardPay shard.
- Cell-Based and Regional IsolationShard-by-customer cohorts for independent failure and scale — ShardPay's cell architecture.
- Multi-Region Active-Active DesignWrite in multiple regions without double-spending — conflict resolution and data residency.
- Conway's Law and Team TopologyArchitecture follows communication paths — align ShardPay service boundaries with team ownership.
- The Distributed Systems Interview PlaybookFramework for trade-off questions, failure scenarios, and architect-level system reasoning.