1. 1.1 Distributed Systemshistorical

    1. What Is a Distributed System - A system composed of multiple subsystems. The subsystems communicate through a network, and each subsystem consists of multiple machines. 2. Why Distributed Systems Are Needed - A single machine has limited read/write capability and is not safe. 3. How to Design a Distributed System 3.1. Replication. 3.2. Partitioning. 4. Theoretical Foundations of Distributed Systems - CAP, BASE.

  2. 1.2 How to Implement Distributed Lockshistorical

    1. What Is a Distributed Lock - A lock in a distributed environment (across processes or machines). It satisfies the following conditions: atomicity, mutual exclusion, no deadlock, and locking and unlocking must be performed by the same client. 2. Why Distributed Locks Are Needed - Built-in locks in various languages, such as Java synchronized and Go mutex, can only guarantee lock properties within a single process and cannot span processes or machines.

  3. 1.3 How to Implement Distributed IDshistorical

    1. What Is a Distributed ID - A unique ID in a distributed environment (multiple machines). Characteristics of distributed IDs: globally unique; roughly increasing; high concurrency; high availability. 2. How to Implement Distributed IDs - Database auto-increment, UUID, Redis-generated IDs, Snowflake IDs.

  4. 1.4 How to Implement Distributed Sessionshistorical

    1. What Is a Distributed Session - A session is server-side memory. - A session shared in a distributed environment (multiple machines). 2. Why Distributed Sessions Are Needed - In a distributed environment, if the traditional Tomcat session mechanism is still used, the following can happen: after a user logs in to system A, the user needs to jump to system B to perform some operations, but the user's login information is stored in system A and system B has no related session.

  5. 1.5 How to Implement Distributed Storagehistorical

    1. What Is Distributed Storage - File storage in a distributed environment (across multiple machines). 2. Why Distributed Storage Is Needed - A single machine has limited performance and capacity. - A single machine does not have high availability. 3. How to Implement Distributed Storage - Distributed-System Replication - Distributed-System Partitioning - Distributed Consistency - Distributed-System Cluster Metadata Management.

  6. 1.6 BASEhistorical

    1. What Is BASE - An extension of AP theory. - It guarantees eventual consistency rather than strong consistency. That is, because failures are unavoidable, I allow data to be different during this period, but after this period the data needs to be consistent. - Availability is obtained by sacrificing strong consistency. When a failure occurs, partial unavailability is allowed, but core functions must remain available.

  7. 1.7 CAPhistorical

    1. Why CAP Exists - A distributed system has multiple nodes, and state needs to be synchronized among the nodes. This requires support from CAP theory. 2. What Is CAP - Only two of the three can be chosen. 2.1. C (Consistency) - Consistency. - A read after a write must return that value. When data is distributed across multiple nodes, the data read from any node must be that value.

  8. 1.8 Distributed-System Cluster Metadata Managementhistorical

    1. What Is Cluster Metadata - Data in any file system can be divided into actual data and metadata. - Data refers to the data we store in files. - Metadata refers to characteristics of files, such as access permissions and data-block distribution. 2. Why Cluster Metadata Is Needed - Only through cluster metadata can we know the mapping relationship between data and partitions.

  9. 1.9 Distributed Consistencyhistorical

    1. What Is Distributed Consistency - Data remains consistent across multiple replicas, that is, data consistency. 2. Why Distributed Consistency Is Needed - In a distributed environment, the same data needs to be replicated to multiple nodes for fault tolerance. Because replication is delayed by network problems, the same data may differ across multiple nodes at the same moment. Distributed consistency exists to solve this problem.

  10. 1.10 Distributed Computinghistorical

    1. What Is Distributed Computing 2. Why Distributed Computing Is Needed 3. Categories of Distributed Computing 4. Batch Processing 5. Stream Processing - Stream processing - Stream: data that gradually increases over time - An event is the smallest unit of stream processing - Each event contains a timestamp indicating its creation time - Events are produced by producers and correspond to multiple consumers.

  11. 1.11 Distributed-System Communicationhistorical

    1. What Is Distributed-System Communication - After a monolithic system is split into multiple subsystems, the subsystems need to communicate to provide a complete service. 2. Communication Methods Between Distributed Systems 2.1. Synchronous Communication - A calls B and waits for the result to return. 2.1.1. REST vs RPC.

  12. 1.12 Distributed-System Service Statehistorical

    1. What Is State - State refers to a program's contextual information or data. - Stateful: a user's previous request and next request are associated, and the request must reach a particular node to be processed normally. - Stateless: a user's previous request and next request are associated, and the request can reach any node and still be processed normally. 2. Stateful Services vs Stateless Services.

  13. 1.13 Distributed-System Upgrade and Rollbackhistorical

    1. What Are Service Upgrade and Rollback - A service in a distributed system has multiple instances. When a new feature is released or an old system is refactored, all of these instances need to be upgraded and replaced. If a problem occurs, rollback needs to happen promptly. 2. Deployment Strategies 2.1. Downtime Deployment - Stop the existing version of the service and then deploy the new version.

  14. 1.14 Distributed-System Failureshistorical

    1. What Is a Distributed-System Failure - A node goes down or the network is unreachable. 2. How to Detect Distributed-System Failures 2.1. Heartbeat Detection 2.2. Gossip Protocol Detection. 3. How to Handle Distributed-System Failures 3.1. Fail-Over - A calls B. If B fails and B has other replicas, A switches to another replica of B.

  15. 1.15 Distributed-System Node Communicationhistorical

    1. What Is Node Communication - Nodes in a distributed system need to communicate and exchange metadata. 1.1. What Is Metadata - The relationship between master and slave, data distribution, and so on. 1.2. Metadata Maintenance Methods - There are generally two ways to maintain metadata in a cluster: centralized and P2P.

  16. 2. Distributed Transactionshistorical

    1. What Are Distributed Transactions - Transactions in a distributed environment. - Traditional transactions operate on a single database node, while distributed transactions span multiple database nodes, or even other data sources such as Redis and MongoDB. 2. Why Distributed Transactions Are Needed - Traditional database transactions can only guarantee transactions within a single database and cannot do so across multiple databases.

  17. 2.1 Distributed Transaction Solution: 2PChistorical

    1. What Is 2PC - A distributed transaction solution that guarantees strong consistency. 2. 2PC Process - Divide the transaction into two phases. - Phase 1: the transaction manager sends prepare requests to all databases. - Phase 2 depends on the result of Phase 1.

  18. 2.2 Distributed Transaction Solution: TCChistorical

    1. What Is TCC - A distributed transaction solution that guarantees eventual consistency. 2. TCC Process - TCC: Try, Confirm, Cancel. - Divide the transaction into two phases. - In Phase 1, execute Try to perform business checks and reserve resources. - Phase 2 depends on the result of Phase 1.

  19. 2.3 Distributed Transaction Solution: Reliable-Message Eventual Consistencyhistorical

    1. What Is Reliable-Message Eventual Consistency - A distributed transaction solution that guarantees eventual consistency. 2. Reliable-Message Eventual-Consistency Process. 3. Use Cases - Suitable for businesses that require eventual consistency, have low time sensitivity, and where the distributed transaction can only succeed rather than fail, such as awarding points after registration or coupons after login.

  20. 2.4 Distributed Transaction Solution: Best-Effort Notificationhistorical

    1. What Is Best-Effort Notification - A distributed transaction solution that guarantees eventual consistency. 2. Best-Effort Notification Process. 3. Use Cases - Suitable for businesses that require eventual consistency, have low time sensitivity, allow a small number of distributed transactions to fail, and where the passive side's processing result does not affect the active side's processing result, such as bank notifications and payment-result notifications.

  21. 2.5 Distributed Transaction Solution: Sagahistorical

    1. What Is Saga - A distributed transaction solution that guarantees eventual consistency. 2. Saga Process - There are multiple transaction participants, and each participant has two pieces of logic: a forward operation and a reverse operation. - Divide the transaction into two phases.

  22. 2.6 Distributed Transaction Solution: 3PChistorical

    1. What Is 3PC - A distributed transaction solution that guarantees strong consistency and is an improved version of 2PC. 2. 3PC Process - Split the first phase of 2PC into two steps, so the whole transaction process has three phases: CanCommit, PreCommit, and DoCommit.

  23. 2.7 Distributed Transaction Solution: Two Phaseshistorical

    1. What Are Two Phases - Divide the entire transaction process into two phases: preparation phase; commit phase or rollback phase. 2. Two-Phase Implementations 2.1. Database Layer - Two-Phase Implementation: 2PC. 2.2. Application Layer - Two-Phase Implementation: TCC.

  24. 2.8 Reliable-Message Eventual-Consistency Implementation: Local Message Tablehistorical

    1. What Is a Local Message Table - Use two transactions, placing an order and then adding points, as an example. The ordering service is Service A and the points service is Service B. 1. Service A executes ordering logic. 2. After Service A successfully places the order, it sends a message to MQ. 3. Service B consumes the message and processes its local transaction.

  25. 2.9 Two-Phase Implementation: 2PChistorical

    1. What Is 2PC - 2PC: Two-Phase Commit. - Divide the transaction into two phases. - In the first phase, the transaction manager sends prepare requests to all databases. - If all respond ok, execute commit in the second phase; if any responds fail, execute rollback in the second phase.

  26. 2.10 Two-Phase Implementation: TCChistorical

    1. What Is TCC - 2PC is a two-phase approach at the database layer, while TCC is a two-phase approach at the application layer. - TCC: Try: attempt to execute the transaction; Confirm: confirm execution of the transaction; Cancel: cancel execution of the transaction. - In essence, it is also a two-phase transaction.

  27. 2.11 Best-Effort Notification Implementation: MQhistorical

    1. What It Is - 1. The producer finishes executing its local transaction and sends a message to MQ. 2. MQ sends the message to the consumer. 3. The consumer consumes the message and executes its local transaction; if successful it ack's, and if it fails it nack's and requeues the message. 4. The consumer can actively call the producer's interface to query message status.

  28. 2.12 RocketMQ Transaction Messageshistorical

    1. What Are RocketMQ Transaction Messages - Traditional local message tables depend on a message table in the database. - RocketMQ transactions encapsulate the local-message-table approach by moving the local message table into MQ, solving the atomicity problem between Producer-side message sending and local transaction execution.

  29. 3. Distributed Consensus Algorithmshistorical

    1. What Are Distributed Consensus Algorithms - More precisely, they are consensus algorithms: making all nodes agree on something. 2. Why Consistency Problems Occur 2.1. Concurrent Client Requests - For example, in a Leader-Follower scenario: one Client, three Nodes A, B, and C. The Client asks A to write x as 1. If A considers x to be 1, then B and C must also consider the value to be 1.

  30. 3.1 Distributed Consensus Algorithm: Paxoshistorical

    1. Basic Paxos 1.1. What Is Basic Paxos - Abbreviated as Paxos. - A distributed consensus algorithm invented by Lamport and the foundation of Raft and ZAB. 1.2. Basic Paxos Algorithm Process 1.2.1. Roles - client: request initiator; not important here. - proposer: proposal proposer, similar to a coordinator.

  31. 3.2 Distributed Consensus Algorithm: ZABhistorical

    - ZAB Protocol.md

  32. 3.3 Distributed Consensus Algorithm: Rafthistorical

    1. What Is Raft - A distributed consensus algorithm invented by Diego Ongaro. 2. Why Raft Is Needed - To solve the complexity of implementing Paxos. 3. Raft Algorithm Process 3.1. Roles - Leader - Follower - Candidate 3.2. Three Phases 3.2.1. Phase 1: Leader Election

  33. 3.4 Distributed Consensus Algorithm: Gossiphistorical

    1. What Is Gossip - An algorithm proposed by Xerox for replicating data among multiple nodes in a distributed database. - Nodes continuously exchange information, and after a period of time all nodes in the cluster will know the complete information. 2. Why Gossip Is Needed 3. Gossip Algorithm Process - Each node periodically and randomly selects a connected node to spread messages.

  34. 3.5 Distributed Consistency Modelshistorical

    1. What Are Distributed Consistency Models - Different consistency models solve consistency problems to different degrees. 2. Categories of Distributed Consistency Models 2.1. Strong Consistency - C in CAP.md - Also called linearizability. 2.2. Weak Consistency - Eventual consistency - Causal consistency - Read-your-writes consistency - Session consistency - Monotonic-read consistency - Monotonic-write consistency - Prefix-read consistency

  35. 4. Distributed-System Replicationhistorical

    1. What Is Replication - The same data is stored on multiple machines. - Each node that stores the data is called a replica. 2. Why Replication Is Needed - Improve availability through data redundancy. - Improve read throughput through read/write separation. 3. Replication Architecture. 4. Replication Methods. 5. Replication Log Formats.

  36. 4.1 Distributed-System Replication Architecture: Leader-Leader Replicationhistorical

    1. What Is Leader-Leader - There are multiple Leaders, and each Leader has multiple Followers. 2. Leader-Leader Use Cases - Multiple data centers. - Applications still need to continue working after the network is disconnected. 3. How Leader-Leader Works 3.1. Leader Election 3.2. Leaders Synchronize Data to Leaders.

  37. 4.2 Distributed-System Replication Architecture: Leaderless Replicationhistorical

    1. What Is Leaderless Replication - There is no Leader. When the client writes, it sends the write request to all replicas in parallel; when it reads, it similarly sends the read request to all replicas in parallel. 2. Leaderless Replication Use Cases. 3. How Leaderless Replication Works 3.1. Data Synchronization 3.1.1. Write-Conflict Problem.

  38. 4.3 Distributed-System Replication Architecture: Leader-Follower Replicationhistorical

    1. What Is Leader-Follower - There is exactly one Leader among the replicas, and all others are Followers. 2. Leader-Follower Use Cases - A single data center. 3. Leader Election - Select one replica as the Leader and use the other replicas as Followers.

  39. 4.4 Distributed-System Replication Logshistorical

    1. What They Are - Data changes between replicas are generally tracked through replication logs, with several formats. 2. Categories 2.1. Physical Logs - Which page was modified, what was the original value, and what is the updated value. 2.2. Logical Logs - Which record was modified, further divided into Statement and Row.

  40. 4.5 Distributed-System Replication Methodshistorical

    1. Synchronous Replication - The Leader succeeds only after synchronizing to all Followers. 1. The client requests the Leader. 2. The Leader writes local data. 3. The Leader synchronizes to Followers. 4. The Leader returns success to the client. 2. Asynchronous Replication - The Leader succeeds once it writes locally.

  41. 4.6 Distributed-System Replication Architecturehistorical

    1. What It Is - Replicas generally have two roles: Leader, responsible for handling client write requests; Follower, synchronizes data from the Leader and can handle client read requests. 2. Categories 2.1. Leader-Follower. 2.2. Leader-Leader. 2.3. Leaderless.

  42. 5. Distributed-System Partitioninghistorical

    1. What Is Partitioning - Split one set of data into multiple parts and store them on different nodes. - There are two layers of mapping: take a field from the data as the key, then map key -> partition; then map partition -> machine/node. 2. Why Partitioning Is Needed - The data volume is too large to store on one node and must be distributed.

  43. 5.1 Distributed-System Partitioning: Data Splittinghistorical

    Distribute keys across partitions/nodes as evenly as possible. 1. Choosing the Split Key 1.1. Primary-Key ID - For example, an auto-increment primary key. The advantage is even data distribution; the disadvantage is that queries by business fields may be slow because all partitions need to be read. 1.2. Business ID - Such as user ID or product ID.

  44. 5.2 Distributed-System Partitioning: Request Processinghistorical

    The routing component routes client read and write requests to the node that contains the corresponding partition. 1. Routing Component. 2. Request Processing 2.1. Add Data - 1. The client generates data containing a sharding key. 2. The client sends the add request to the routing component.

  45. 5.3 Distributed-System Partitioning: Routing Componentshistorical

    1. client - The client locally stores the relationship between partitions and server nodes and directly requests the correct node. 2. proxy - The client requests the routing layer, and the routing layer forwards the request to the correct node. 3. server - The client requests any server node, and that server node forwards the request to the correct node.

  46. 5.4 Distributed-System Partitioning: Partition Assignmenthistorical

    Distribute partitions/nodes across machines as evenly as possible. 1. Assignment Methods 1.1. Static Assignment - Create far more nodes than machines. Advantage: when migrating nodes to other machines, the cluster can still respond externally. Disadvantage: the maximum number of machines is fixed.