Distributed Systems Topics Covered in PKUFlyingPig/cs-self-learning
The Distributed Systems section in this repository delivers a comprehensive, architecture-centric curriculum spanning MIT 6.824 fundamentals—Raft consensus, fault-tolerant storage, and CAP theorem—through modern distributed ML training stacks like PyTorch DDP and ZeRO.
The PKUFlyingPig/cs-self-learning repository structures its Distributed Systems content around the classic MIT 6.824 course, augmented with curated resources for large-scale machine learning systems. This learning path moves semantically from theoretical consensus algorithms to hands-on implementation of replicated state machines, culminating in production-grade distributed training architectures.
Core MIT 6.824 Curriculum
The backbone of the distributed systems education resides in docs/并行与分布式系统/MIT6.824.md (Chinese) and docs/Parallel and Distributed Systems/MIT6.824.en.md (English), which organize content around seminal research papers and their practical implementations.
Fundamentals and Classic Papers
Each lecture centers on a landmark paper that defined the field. You will study Paxos, Raft, Chord, and Spanner, with notes explaining the specific problem each solves, the underlying algorithmic ideas, and the engineering trade-offs involved. This paper-centric approach ensures you understand not just how distributed systems work, but why they are designed that way.
Consensus and State Machine Replication
The curriculum provides in-depth coverage of the Raft consensus algorithm, including leader election, log replication, and the safety and liveness guarantees required for correct state machine replication. As implemented in the course projects, you build a fault-tolerant key-value store where the Raft struct manages distributed state:
type Raft struct {
mu sync.Mutex
peers []string // other server addresses
me int // this server's index
state ServerState
currentTerm int
votedFor *int
log []LogEntry
commitIndex int
lastApplied int
}
// RequestVote RPC handler (simplified)
func (rf *Raft) RequestVote(args *RequestVoteArgs, reply *RequestVoteReply) error {
rf.mu.Lock()
defer rf.mu.Unlock()
if args.Term < rf.currentTerm {
reply.VoteGranted = false
return nil
}
// vote if we haven't voted yet and candidate's log is at least as up-to-date
if rf.votedFor == nil || *rf.votedFor == args.CandidateId {
if args.LastLogTerm > rf.lastLogTerm() ||
(args.LastLogTerm == rf.lastLogTerm() && args.LastLogIndex >= rf.lastLogIndex()) {
rf.votedFor = &args.CandidateId
rf.currentTerm = args.Term
reply.VoteGranted = true
return nil
}
}
reply.VoteGranted = false
return nil
}
This RequestVote implementation demonstrates the term-based logic and log-comparison checks required for safe leader election in a replicated log system.
Fault Tolerance and Crash Recovery
The materials cover techniques for handling node failures, including checkpointing and snapshotting mechanisms that allow services to recover state after crashes. You learn to build robust KV-store services that maintain availability despite individual node failures, implementing the full stack from RPC handling to persistent state management across the four course labs.
Consistency Models and Storage Architecture
CAP Theorem and Consistency Guarantees
The repository explores the CAP theorem and the fundamental trade-offs between consistency, availability, and partition tolerance. You will analyze formal consistency models including linearizability, sequential consistency, and eventual consistency, understanding which guarantees are necessary for specific application requirements.
Distributed Storage and KV-Stores
Practical system design focuses on building a sharded, fault-tolerant key-value store atop Raft. Topics include sharding strategies, partitioning schemes, and client-side routing mechanisms that allow the system to scale horizontally while maintaining strong consistency semantics.
Scalability and Load Balancing
The curriculum addresses horizontal scaling through consistent hashing, request forwarding patterns, and load-balancing proxy architectures. These techniques enable distributed systems to handle increasing load by adding nodes rather than upgrading hardware.
Distributed Machine Learning Systems
Extending beyond traditional distributed systems, the repository includes docs/深度生成模型/大语言模型/CMU11-868.en.md (referenced in docs/机器学习系统/CMU11-868.en.md), which treats distributed training as a systems problem. This section covers modern stacks including DDP (Distributed Data Parallel), GPipe, Megatron-LM, and ZeRO, showing how consensus and fault-tolerance concepts apply to large-scale ML workloads.
The implementation uses PyTorch's distributed training APIs to demonstrate gradient synchronization across nodes:
import os
import torch
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.utils.data import DataLoader, DistributedSampler
from torchvision import datasets, transforms
def setup(rank, world_size):
os.environ["MASTER_ADDR"] = "localhost"
os.environ["MASTER_PORT"] = "12355"
dist.init_process_group("nccl", rank=rank, world_size=world_size)
def cleanup():
dist.destroy_process_group()
def main(rank, world_size):
setup(rank, world_size)
# model + optimizer
model = torchvision.models.resnet50().to(rank)
ddp_model = DDP(model, device_ids=[rank])
# Distributed sampler ensures each GPU sees a disjoint subset
dataset = datasets.CIFAR10(root="./data", download=True,
transform=transforms.ToTensor())
sampler = DistributedSampler(dataset, num_replicas=world_size, rank=rank)
loader = DataLoader(dataset, batch_size=128, sampler=sampler)
loss_fn = torch.nn.CrossEntropyLoss()
optimizer = torch.optim.SGD(ddp_model.parameters(), lr=0.01)
for epoch in range(10):
sampler.set_epoch(epoch)
for xb, yb in loader:
optimizer.zero_grad()
preds = ddp_model(xb.to(rank))
loss = loss_fn(preds, yb.to(rank))
loss.backward()
optimizer.step()
cleanup()
# launch with: torchrun --nproc_per_node=4 your_script.py
This snippet illustrates the same coordination principles underlying Raft—processes must agree on a common optimization step—applied to gradient synchronization across GPUs.
Debugging and Testing Distributed Code
The MIT 6.824 project guidance in docs/CS学习规划.en.md emphasizes deterministic replay and techniques for reproducing race conditions. You learn to write unit and integration tests that verify correctness under network partitions and node failures, using tools that simulate unreliable networks to expose concurrency bugs before deployment.
Supplemental Resources
The docs/好书推荐.md file curates essential external reading including Patterns of Distributed Systems and Distributed Systems for Fun and Profit, providing practical engineering perspectives that complement the academic rigor of the MIT coursework.
Summary
- MIT 6.824 forms the pedagogical core, with Chinese and English documentation in
docs/Parallel and Distributed Systems/MIT6.824.en.mdcovering Raft, Paxos, and Spanner. - Consensus algorithms are taught through implementation, including the
RequestVoteRPC logic and state machine replication patterns. - Fault tolerance mechanisms include checkpointing, snapshotting, and crash-recovery protocols for KV-stores.
- CAP theorem and consistency models (linearizability, sequential, eventual) define the theoretical boundaries of system design.
- Distributed ML training extends the curriculum to modern GPU clusters using PyTorch DDP, as detailed in the CMU 11-868 materials.
- Testing methodologies emphasize deterministic replay and race-condition detection for concurrent distributed code.
Frequently Asked Questions
What prerequisites are needed for the Distributed Systems track?
You should have solid systems programming experience in Go (for MIT 6.824 labs) or C/C++, plus understanding of operating systems and networking fundamentals. The repository assumes familiarity with concurrent programming and basic RPC concepts before tackling the Raft implementation.
How does the repository cover modern distributed training alongside classic consensus theory?
The curriculum bridges theory and practice by first establishing consensus fundamentals (Raft, Paxos) in the MIT 6.824 section, then applying these coordination concepts to distributed ML systems in docs/深度生成模型/大语言模型/CMU11-868.en.md. This shows how gradient synchronization in PyTorch DDP requires the same process-agreement mechanisms underlying traditional state machine replication.
Where can I find the implementation details for the Raft consensus project?
Complete architectural notes and lab specifications reside in docs/并行与分布式系统/MIT6.824.md (Chinese) and docs/Parallel and Distributed Systems/MIT6.824.en.md (English). These files detail the Raft struct definition, RPC handlers like RequestVote, and the four-lab progression from basic consensus to a sharded fault-tolerant KV-store.
What additional reading does the repository recommend for distributed systems engineering?
The docs/好书推荐.md file lists curated blogs including Patterns of Distributed Systems and Distributed Systems for Fun and Profit, which provide practical engineering patterns and real-world debugging strategies that complement the academic paper focus of the main curriculum.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →