VibeKoding / Ensiklopedia ยท Fondasi KuatEnsiklopedia ยท Fondasi Kuat / Principles of Distributed SystemsPrinciples of Distributed Systems
VK

Principles of Distributed SystemsPrinciples of Distributed Systems

๐Ÿ“š Ensiklopedia ยท Fondasi KuatEnsiklopedia ยท Fondasi Kuat ๐ŸŒ Dual Bahasa (ID / EN) โšก VibeKoding Native

Ensiklopedia VibeKoding: Principles of Distributed Systems.Ensiklopedia VibeKoding: Principles of Distributed Systems.

๐Ÿ’ก Tips Praktis๐Ÿ’ก Pro Tip

When one machine is no longer enough, the real problems begin. Distributed systems are the backbone of the modern internet โ€” from WeChat messaging to Taobao orders, hundreds or thousands of machines work together behind the scenes. But "distributed" is no free lunch; it introduces a host of challenges that single-machine systems never face.When one machine is no longer enough, the real problems begin. Distributed systems are the backbone of the modern internet โ€” from WeChat messaging to Taobao orders, hundreds or thousands of machines work together behind the scenes. But "distributed" is no free lunch; it introduces a host of challenges that single-machine systems never face.

What will you learn from this article?What will you learn from this article?

After reading this chapter, you will gain:After reading this chapter, you will gain:

ChapterContentCore Concepts
Chapter 1Why Distributed SystemsScalability, availability, geographic distribution
Chapter 2CAP TheoremConsistency, availability, partition tolerance
Chapter 3Consistency ModelsStrong consistency, eventual consistency, causal consistency
Chapter 4Eight ChallengesNetwork, clocks, partitions, split brain, etc.
Chapter 5Consensus AlgorithmsPaxos, Raft, ZAB
Chapter 6Distributed Transactions2PC, Saga, TCC

------

0. The Big Picture: Motivation for needing Distributed Systems0. The Big Picture: Motivation for needing Distributed Systems

Single-machine systems are simple and reliable, but they have three insurmountable bottlenecks:Single-machine systems are simple and reliable, but they have three insurmountable bottlenecks:

BottleneckDescriptionDistributed Solution
Performance ceilingA single machine has physical limits on CPU, memory, and diskHorizontal scaling: add more machines to share the load
Single point of failureIf one machine goes down, the entire service goes downRedundant replicas: multiple machines serve as backups for each other
Geographic latencyUsers are spread across the globe; a single machine can only be in one placeMulti-region deployment: serve users from nearby locations
๐Ÿ’ก Tips Praktis๐Ÿ’ก Pro Tip

Distributed systems solve the problems above but introduce new complexities: unreliable networks, clock skew, partial failures, data consistency... These are the "challenges" this article will discuss. Peter Deutsch's Eight Fallacies of Distributed Computing tell us that the following assumptions are all wrong in distributed environments: 1. The network is reliable 2. Latency is zero 3. Bandwidth is infinite 4. The network is secure 5. Topology doesn't change 6. There is one administrator 7. Transport cost is zero 8. The network is homogeneousDistributed systems solve the problems above but introduce new complexities: unreliable networks, clock skew, partial failures, data consistency... These are the "challenges" this article will discuss. Peter Deutsch's Eight Fallacies of Distributed Computing tell us that the following assumptions are all wrong in distributed environments: 1. The network is reliable 2. Latency is zero 3. Bandwidth is infinite 4. The network is secure 5. Topology doesn't change 6. There is one administrator 7. Transport cost is zero 8. The network is homogeneous

------

1. CAP Theorem: The "Impossible Triangle" of Distributed Systems1. CAP Theorem: The "Impossible Triangle" of Distributed Systems

In 2000, Eric Brewer proposed the CAP conjecture (later proven as a theorem): a distributed system can satisfy at most two of the following three properties simultaneously.In 2000, Eric Brewer proposed the CAP conjecture (later proven as a theorem): a distributed system can satisfy at most two of the following three properties simultaneously.

PropertyMeaningIntuitive Explanation
ConsistencyAll nodes see the same data at the same timeYou check your balance at any ATM and get the same result
AvailabilityEvery request receives a non-error responseThe system always responds to you; it never says "service unavailable"
Partition toleranceThe system continues to operate during a network partitionEven if some network cables are cut, the system still works

Motivation for Onlying Choose TwoMotivation for Onlying Choose Two

In a distributed environment, network partitions (P) are inevitable โ€” fiber optic cables get dug up, switches fail, data centers lose connectivity. So P is mandatory, and the real choice is a trade-off between C and A:In a distributed environment, network partitions (P) are inevitable โ€” fiber optic cables get dug up, switches fail, data centers lose connectivity. So P is mandatory, and the real choice is a trade-off between C and A:

๐Ÿ’ก Tips Praktis๐Ÿ’ก Pro Tip

Real-world systems are not simply "CP or AP." Many systems make different choices for different operations โ€” for example, in the same database, reads can be AP (allowing stale reads) while writes are CP (requiring majority acknowledgment).Real-world systems are not simply "CP or AP." Many systems make different choices for different operations โ€” for example, in the same database, reads can be AP (allowing stale reads) while writes are CP (requiring majority acknowledgment).

------

2. Consistency Models: The "Strictness" of Data Synchronization2. Consistency Models: The "Strictness" of Data Synchronization

Consistency is not a binary switch (on or off); it is a spectrum. Different consistency models make different trade-offs between "correctness" and "performance."Consistency is not a binary switch (on or off); it is a spectrum. Different consistency models make different trade-offs between "correctness" and "performance."

Consistency Model ComparisonConsistency Model Comparison

ModelGuaranteeLatencyUse Cases
Strong consistencyReads always return the most recently written valueHigh (requires waiting for sync)Bank transfers, inventory deduction
Eventual consistencyAll replicas will eventually converge, but intermediate reads may be staleLow (writes return immediately)Social feeds, DNS
Causal consistencyCausally related operations are guaranteed to be orderedMediumComment replies, collaborative editing
LinearizabilityAll operations appear to execute sequentially as if on a single machineHighestDistributed locks, leader election
Session consistencyWithin the same session, reads reflect your own writesLow-MediumUser personal data
๐Ÿ’ก Tips Praktis๐Ÿ’ก Pro Tip

The most common practical requirement is: after a user modifies their data, they can immediately see the update (but other users may see it later). This is called "Read Your Own Writes" consistency, a practical enhancement over eventual consistency.The most common practical requirement is: after a user modifies their data, they can immediately see the update (but other users may see it later). This is called "Read Your Own Writes" consistency, a practical enhancement over eventual consistency.

------

3. Eight Challenges: The "Minefield" of Distributed Systems3. Eight Challenges: The "Minefield" of Distributed Systems

The complexity of distributed systems doesn't come from any single problem but from multiple problems intertwining. Here are the eight core challenges.The complexity of distributed systems doesn't come from any single problem but from multiple problems intertwining. Here are the eight core challenges.

Relationships Between ChallengesRelationships Between Challenges

These eight challenges are not isolated; they are interconnected:These eight challenges are not isolated; they are interconnected:

๐Ÿ’ก Tips Praktis๐Ÿ’ก Pro Tip

There is no "perfect" solution in distributed systems, only "appropriate" trade-offs. Understanding the nature of these challenges is essential for making the right design decisions.There is no "perfect" solution in distributed systems, only "appropriate" trade-offs. Understanding the nature of these challenges is essential for making the right design decisions.

------

4. Consensus Algorithms: How to Get Multiple Machines to "Agree"4. Consensus Algorithms: How to Get Multiple Machines to "Agree"

Consensus algorithms are at the heart of distributed systems โ€” they solve the problem of how multiple nodes can agree on a value, even when some nodes fail or the network is slow.Consensus algorithms are at the heart of distributed systems โ€” they solve the problem of how multiple nodes can agree on a value, even when some nodes fail or the network is slow.

4.1 Paxos4.1 Paxos

Proposed by Leslie Lamport in 1990, it was the first consensus algorithm to be rigorously proven correct.Proposed by Leslie Lamport in 1990, it was the first consensus algorithm to be rigorously proven correct.

RoleResponsibility
ProposerProposes a value
AcceptorVotes to accept or reject proposals
LearnerLearns the final chosen value

Two-Phase Process:Two-Phase Process:

  1. Prepare phase: The Proposer sends a proposal number; Acceptors promise not to accept proposals with smaller numbersPrepare phase: The Proposer sends a proposal number; Acceptors promise not to accept proposals with smaller numbers
  2. Accept phase: The Proposer sends the actual value; if a majority of Acceptors accept, the proposal passesAccept phase: The Proposer sends the actual value; if a majority of Acceptors accept, the proposal passes
  3. ๐Ÿ’ก Tips Praktis๐Ÿ’ก Pro Tip

    Paxos is correct but notoriously difficult to understand and implement. Lamport's own paper used a Greek parliament analogy, which ended up confusing even more people.Paxos is correct but notoriously difficult to understand and implement. Lamport's own paper used a Greek parliament analogy, which ended up confusing even more people.

    4.2 Raft: Built for Understandability4.2 Raft: Built for Understandability

    In 2014, Diego Ongaro proposed Raft with the goal of creating "an understandable Paxos." It decomposes consensus into three sub-problems:In 2014, Diego Ongaro proposed Raft with the goal of creating "an understandable Paxos." It decomposes consensus into three sub-problems:

    Sub-problemDescription
    Leader electionElect a Leader in the cluster; all writes go through the Leader
    Log replicationThe Leader replicates operation logs to all Followers
    SafetyGuarantees that committed logs will never be overwritten

    Raft's Core Process:Raft's Core Process:

    1. When the cluster starts, all nodes are FollowersWhen the cluster starts, all nodes are Followers
    2. If a Follower times out without receiving a Leader heartbeat, it becomes a Candidate and initiates an electionIf a Follower times out without receiving a Leader heartbeat, it becomes a Candidate and initiates an election
    3. The Candidate that receives a majority of votes becomes the new LeaderThe Candidate that receives a majority of votes becomes the new Leader
    4. The Leader accepts client requests and commits logs after replicating them to a majority of nodesThe Leader accepts client requests and commits logs after replicating them to a majority of nodes
    5. 4.3 Consensus Algorithm Comparison4.3 Consensus Algorithm Comparison

      AlgorithmYear ProposedUnderstandabilitySystems Using It
      Paxos1990DifficultGoogle Chubby
      Raft2014Easyetcd, Consul, TiKV
      ZAB2011MediumZooKeeper
      EPaxos2013DifficultPrimarily academic research

      ------

      5. Distributed Transactions: Cross-Node "All-or-Nothing"5. Distributed Transactions: Cross-Node "All-or-Nothing"

      Single-machine database transactions achieve ACID through local locks and logs. But when a business operation involves multiple services or databases, how do you ensure atomicity?Single-machine database transactions achieve ACID through local locks and logs. But when a business operation involves multiple services or databases, how do you ensure atomicity?

      5.1 Two-Phase Commit (2PC)5.1 Two-Phase Commit (2PC)

      The most classic distributed transaction protocol, divided into two phases:The most classic distributed transaction protocol, divided into two phases:

      PhaseCoordinator ActionParticipant Action
      PrepareAsks all participants "Can you commit?"Executes the operation but does not commit; replies Yes/No
      CommitIf all Yes, sends CommitFormally commits; if any No, all roll back

      Problems with 2PC:Problems with 2PC:

      • Blocking: After Prepare, if the coordinator goes down, participants will wait indefinitelyBlocking: After Prepare, if the coordinator goes down, participants will wait indefinitely
      • Single point of failure: The coordinator is a single point; if it fails, the entire transaction stallsSingle point of failure: The coordinator is a single point; if it fails, the entire transaction stalls
      • Poor performance: Requires multiple network round trips and holds locks for a long timePoor performance: Requires multiple network round trips and holds locks for a long time

      5.2 Saga Pattern5.2 Saga Pattern

      Saga breaks a large transaction into multiple local transactions, each with a corresponding compensating action. If any step fails, compensations are executed in reverse order.Saga breaks a large transaction into multiple local transactions, each with a corresponding compensating action. If any step fails, compensations are executed in reverse order.

      E-commerce Order Saga Example:E-commerce Order Saga Example:

      StepForward OperationCompensating Operation
      T1Create order (pending payment)Cancel order
      T2Deduct inventoryRestore inventory
      T3Deduct balanceRefund balance
      T4Confirm order (paid)โ€”

      If T3 (deduct balance) fails: execute C2 (restore inventory) โ†’ C1 (cancel order).If T3 (deduct balance) fails: execute C2 (restore inventory) โ†’ C1 (cancel order).

      Two Orchestration Approaches:Two Orchestration Approaches:

      • Choreography: Each service listens for events and decides its own next step. Simple but hard to track global stateChoreography: Each service listens for events and decides its own next step. Simple but hard to track global state
      • Orchestration: A central coordinator controls the workflow. Clear but the coordinator is a single pointOrchestration: A central coordinator controls the workflow. Clear but the coordinator is a single point

      5.3 TCC (Try-Confirm-Cancel)5.3 TCC (Try-Confirm-Cancel)

      TCC is a business-layer implementation of 2PC, splitting each operation into three phases:TCC is a business-layer implementation of 2PC, splitting each operation into three phases:

      PhaseDescriptionExample (Deduct Inventory)
      TryReserve resources without actually executingFreeze 10 units of inventory (available -10, frozen +10)
      ConfirmConfirm execution, consume reserved resourcesFrozen -10 (actual deduction)
      CancelCancel reservation, release resourcesFrozen -10, available +10 (restore)

      5.4 Comparison of Three Approaches5.4 Comparison of Three Approaches

      ApproachConsistencyPerformanceComplexityUse Cases
      2PCStrong consistencyLowMediumCross-database transactions at the database layer
      SagaEventual consistencyHighHighLong-running business processes (orders, logistics)
      TCCEventual consistencyMediumHighestHigh-reliability financial scenarios
      ๐Ÿ’ก Tips Praktis๐Ÿ’ก Pro Tip

      - If you can use a single-database transaction, don't use distributed transactions - Saga + message queues are sufficient for most business scenarios - TCC is suitable for financial scenarios requiring extremely high consistency, but development costs are high - 2PC is suitable for automatic handling by database middleware (e.g., ShardingSphere)- If you can use a single-database transaction, don't use distributed transactions - Saga + message queues are sufficient for most business scenarios - TCC is suitable for financial scenarios requiring extremely high consistency, but development costs are high - 2PC is suitable for automatic handling by database middleware (e.g., ShardingSphere)

      ------

      SummarySummary

      Distributed systems are the infrastructure of the modern internet, but their complexity far exceeds that of single-machine systems. Understanding these challenges is not about "solving" them (many are fundamental), but about making the right trade-offs when designing systems.Distributed systems are the infrastructure of the modern internet, but their complexity far exceeds that of single-machine systems. Understanding these challenges is not about "solving" them (many are fundamental), but about making the right trade-offs when designing systems.

      Key takeaways from this chapter:Key takeaways from this chapter:

      1. CAP Theorem: Network partitions are inevitable; the real choice is a trade-off between consistency and availabilityCAP Theorem: Network partitions are inevitable; the real choice is a trade-off between consistency and availability
      2. Consistency Models: From strong to eventual consistency is a spectrum; choose based on business requirementsConsistency Models: From strong to eventual consistency is a spectrum; choose based on business requirements
      3. Eight Challenges: Unreliable networks, clock skew, network partitions, split brain, and more are all interconnectedEight Challenges: Unreliable networks, clock skew, network partitions, split brain, and more are all interconnected
      4. Consensus Algorithms: Raft is currently the most practical consensus algorithm; etcd/Consul are built on itConsensus Algorithms: Raft is currently the most practical consensus algorithm; etcd/Consul are built on it
      5. Distributed Transactions: Saga for most scenarios, TCC for financial scenarios, 2PC for the database layerDistributed Transactions: Saga for most scenarios, TCC for financial scenarios, 2PC for the database layer
      6. Further ReadingFurther Reading

        • [Designing Data-Intensive Applications](https://dataintensive.net/) - Martin Kleppmann's distributed systems classic[Designing Data-Intensive Applications](https://dataintensive.net/) - Martin Kleppmann's distributed systems classic
        • [The Raft Consensus Algorithm](https://raft.github.io/) - Official Raft visualization demo[The Raft Consensus Algorithm](https://raft.github.io/) - Official Raft visualization demo
        • [CAP Twelve Years Later](https://www.infoq.com/articles/cap-twelve-years-later-how-the-rules-have-changed/) - Brewer's reassessment of CAP[CAP Twelve Years Later](https://www.infoq.com/articles/cap-twelve-years-later-how-the-rules-have-changed/) - Brewer's reassessment of CAP
        • [Jepsen](https://jepsen.io/) - Distributed systems correctness testing framework[Jepsen](https://jepsen.io/) - Distributed systems correctness testing framework
        • [Distributed Systems Patterns](https://martinfowler.com/articles/patterns-of-distributed-systems/) - Martin Fowler's distributed patterns collection[Distributed Systems Patterns](https://martinfowler.com/articles/patterns-of-distributed-systems/) - Martin Fowler's distributed patterns collection