Ensiklopedia VibeKoding: Principles of Distributed Systems.Ensiklopedia VibeKoding: Principles of Distributed Systems.
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:
| Chapter | Content | Core Concepts |
|---|---|---|
| Chapter 1 | Why Distributed Systems | Scalability, availability, geographic distribution |
| Chapter 2 | CAP Theorem | Consistency, availability, partition tolerance |
| Chapter 3 | Consistency Models | Strong consistency, eventual consistency, causal consistency |
| Chapter 4 | Eight Challenges | Network, clocks, partitions, split brain, etc. |
| Chapter 5 | Consensus Algorithms | Paxos, Raft, ZAB |
| Chapter 6 | Distributed Transactions | 2PC, Saga, TCC |
------
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:
| Bottleneck | Description | Distributed Solution |
|---|---|---|
| Performance ceiling | A single machine has physical limits on CPU, memory, and disk | Horizontal scaling: add more machines to share the load |
| Single point of failure | If one machine goes down, the entire service goes down | Redundant replicas: multiple machines serve as backups for each other |
| Geographic latency | Users are spread across the globe; a single machine can only be in one place | Multi-region deployment: serve users from nearby locations |
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
------
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.
| Property | Meaning | Intuitive Explanation |
|---|---|---|
| Consistency | All nodes see the same data at the same time | You check your balance at any ATM and get the same result |
| Availability | Every request receives a non-error response | The system always responds to you; it never says "service unavailable" |
| Partition tolerance | The system continues to operate during a network partition | Even if some network cables are cut, the system still works |
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:
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).
------
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."
| Model | Guarantee | Latency | Use Cases |
|---|---|---|---|
| Strong consistency | Reads always return the most recently written value | High (requires waiting for sync) | Bank transfers, inventory deduction |
| Eventual consistency | All replicas will eventually converge, but intermediate reads may be stale | Low (writes return immediately) | Social feeds, DNS |
| Causal consistency | Causally related operations are guaranteed to be ordered | Medium | Comment replies, collaborative editing |
| Linearizability | All operations appear to execute sequentially as if on a single machine | Highest | Distributed locks, leader election |
| Session consistency | Within the same session, reads reflect your own writes | Low-Medium | User personal data |
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.
------
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.
These eight challenges are not isolated; they are interconnected:These eight challenges are not isolated; they are interconnected:
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.
------
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.
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.
| Role | Responsibility |
|---|---|
| Proposer | Proposes a value |
| Acceptor | Votes to accept or reject proposals |
| Learner | Learns the final chosen value |
Two-Phase Process:Two-Phase Process:
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.
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-problem | Description |
|---|---|
| Leader election | Elect a Leader in the cluster; all writes go through the Leader |
| Log replication | The Leader replicates operation logs to all Followers |
| Safety | Guarantees that committed logs will never be overwritten |
Raft's Core Process:Raft's Core Process:
| Algorithm | Year Proposed | Understandability | Systems Using It |
|---|---|---|---|
| Paxos | 1990 | Difficult | Google Chubby |
| Raft | 2014 | Easy | etcd, Consul, TiKV |
| ZAB | 2011 | Medium | ZooKeeper |
| EPaxos | 2013 | Difficult | Primarily academic research |
------
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?
The most classic distributed transaction protocol, divided into two phases:The most classic distributed transaction protocol, divided into two phases:
| Phase | Coordinator Action | Participant Action |
|---|---|---|
| Prepare | Asks all participants "Can you commit?" | Executes the operation but does not commit; replies Yes/No |
| Commit | If all Yes, sends Commit | Formally commits; if any No, all roll back |
Problems with 2PC:Problems with 2PC:
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:
| Step | Forward Operation | Compensating Operation |
|---|---|---|
| T1 | Create order (pending payment) | Cancel order |
| T2 | Deduct inventory | Restore inventory |
| T3 | Deduct balance | Refund balance |
| T4 | Confirm 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:
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:
| Phase | Description | Example (Deduct Inventory) |
|---|---|---|
| Try | Reserve resources without actually executing | Freeze 10 units of inventory (available -10, frozen +10) |
| Confirm | Confirm execution, consume reserved resources | Frozen -10 (actual deduction) |
| Cancel | Cancel reservation, release resources | Frozen -10, available +10 (restore) |
| Approach | Consistency | Performance | Complexity | Use Cases |
|---|---|---|---|---|
| 2PC | Strong consistency | Low | Medium | Cross-database transactions at the database layer |
| Saga | Eventual consistency | High | High | Long-running business processes (orders, logistics) |
| TCC | Eventual consistency | Medium | Highest | High-reliability financial scenarios |
- 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)
------
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: