返回文章列表
技术2026年9月18日18 分钟阅读

共识算法:从Paxos到Raft的形式化理解与工程实践

共识算法:从Paxos到Raft的形式化理解与工程实践

一、问题定义:分布式共识的本质

1.1 什么是共识?

在分布式系统中,共识(Consensus) 是指多个节点就某个值达成一致的过程。看似简单的目标,在异步网络、节点故障、消息丢失的环境下,却成为计算机科学中最具挑战性的问题之一。

形式化定义:

给定 nn 个进程组成的集合 P={p1,p2,...,pn}P = \{p_1, p_2, ..., p_n\},每个进程 pip_i 提出一个值 viv_i,共识算法需要满足:

  1. 终止性(Termination):每个正确进程最终决定是否接受某个值
  2. 一致性(Agreement):所有正确进程接受相同的值
  3. 有效性(Validity):被接受的值必须是某个进程提出的值

1.2 系统模型与故障类型

通信模型:

模型 消息延迟 适用场景
同步 有上界 数据中心内部
异步 无保证 广域网、互联网
部分同步 最终有上界 实际分布式系统

故障类型:

故障类型演进:

1. 崩溃停止(Crash-Stop)
   └── 节点停止运行,不再响应
   └── 最简单模型

2. 崩溃恢复(Crash-Recovery)
   └── 节点崩溃后可能恢复
   └── 需要持久化存储

3. 拜占庭故障(Byzantine)
   └── 节点可能发送任意错误消息
   └── 恶意行为或硬件故障
   └── 需要 BFT 算法(PBFT、Tendermint)

1.3 FLP不可能性结果

1985年,Fischer、Lynch 和 Paterson 发表了著名的 FLP不可能性结果:

在异步系统中,即使只有一个进程可能故障,也不存在确定性的共识算法。

证明思路:

核心洞察:在异步系统中,无法区分
1. 进程崩溃
2. 消息延迟极长

因此,算法必须在"等待可能延迟的消息"和"继续推进"之间做出选择。
无论怎么选择,都存在一种执行路径导致无法达成共识。

实际意义:

FLP 结果并不意味着共识不可能实现,而是告诉我们:

  1. 必须使用随机化(打破确定性)
  2. 或引入超时机制(将异步转为同步)
  3. 或容忍活性失效(牺牲终止性保证安全性)

二、Paxos算法:分布式共识的基石

2.1 Lamport的希腊寓言

1990年,Leslie Lamport 用虚构的希腊议会故事描述了 Paxos 算法。这个看似古怪的叙述方式,实际上揭示了算法的核心思想:

即使议员们可能离开、消息可能丢失,议会仍能通过一系列投票程序通过法令。

2.2 Paxos的三种角色

Paxos 角色:

1. Proposer(提议者)
   └── 提出值(Value)
   └── 推动共识达成
   └── 可以有多个

2. Acceptor(接受者)
   └── 对提议进行投票
   └── 存储接受的值
   └── 关键:多数派(Quorum)

3. Learner(学习者)
   └── 学习已达成共识的值
   └── 不参与决策过程

多数派(Quorum):

定义:Acceptor 的子集,满足任意两个 Quorum 有交集

性质:
- 如果有 $N$ 个 Acceptor,Quorum 大小至少为 $\lfloor N/2 \rfloor + 1$
- 任意两个 Quorum 至少共享一个 Acceptor
- 这是保证一致性的关键

2.3 两阶段协议

Paxos 通过两阶段提交达成共识:

Phase 1:Prepare/Promise

Proposer                    Acceptor
   │                           │
   │  1. Prepare(n)            │
   │ ─────────────────────────>│
   │  (提议编号 n)              │
   │                           │
   │  2. Promise(n, v)         │
   │ <─────────────────────────│
   │  (承诺不接受更小编号)       │
   │  (返回已接受的值 v)         │
   │                           │

Phase 2:Accept/Accepted

Proposer                    Acceptor
   │                           │
   │  3. Accept(n, v)          │
   │ ─────────────────────────>│
   │  (提议值 v)                │
   │                           │
   │  4. Accepted(n, v)        │
   │ <─────────────────────────│
   │  (确认接受)                │
   │                           │

2.4 Python实现

python
import random
import threading
import time
from typing import Optional, Set, Dict, List
from dataclasses import dataclass
from enum import Enum

class PaxosState(Enum):
    IDLE = "idle"
    PREPARING = "preparing"
    ACCEPTING = "accepting"
    DECIDED = "decided"

@dataclass
class Proposal:
    number: int
    value: any
    
    def __eq__(self, other):
        return self.number == other.number and self.value == other.value

class Acceptor:
    """
    Paxos Acceptor 实现
    """
    
    def __init__(self, acceptor_id: str):
        self.acceptor_id = acceptor_id
        self.promised_proposal: Optional[int] = None
        self.accepted_proposal: Optional[Proposal] = None
        self.lock = threading.Lock()
    
    def prepare(self, proposal_number: int) -> Optional[Proposal]:
        """
        Phase 1: 处理 Prepare 请求
        
        返回:
        - 如果承诺,返回已接受的值(如果有)
        - 如果拒绝,返回 None
        """
        with self.lock:
            if self.promised_proposal is None or \
               proposal_number > self.promised_proposal:
                self.promised_proposal = proposal_number
                return self.accepted_proposal
            return None  # 拒绝
    
    def accept(self, proposal: Proposal) -> bool:
        """
        Phase 2: 处理 Accept 请求
        
        返回:
        - True: 接受
        - False: 拒绝
        """
        with self.lock:
            if self.promised_proposal is not None and \
               proposal.number >= self.promised_proposal:
                self.accepted_proposal = proposal
                return True
            return False
    
    def get_accepted(self) -> Optional[Proposal]:
        """获取已接受的值"""
        with self.lock:
            return self.accepted_proposal

class Proposer:
    """
    Paxos Proposer 实现
    """
    
    def __init__(self, proposer_id: str, acceptors: List[Acceptor], quorum_size: int):
        self.proposer_id = proposer_id
        self.acceptors = acceptors
        self.quorum_size = quorum_size
        self.proposal_counter = 0
        self.lock = threading.Lock()
    
    def _generate_proposal_number(self) -> int:
        """生成唯一的提议编号"""
        with self.lock:
            self.proposal_counter += 1
            # 使用 (round, proposer_id) 确保唯一性
            return self.proposal_counter * 1000 + hash(self.proposer_id) % 1000
    
    def propose(self, value: any) -> Optional[any]:
        """
        执行完整的 Paxos 两阶段协议
        
        返回:
        - 达成共识的值
        - None 如果失败
        """
        proposal_number = self._generate_proposal_number()
        
        # Phase 1: Prepare
        promises = []
        highest_accepted: Optional[Proposal] = None
        
        for acceptor in self.acceptors:
            accepted = acceptor.prepare(proposal_number)
            if accepted is not None:  # 承诺
                promises.append(acceptor)
                # 记录已接受的最高编号值
                if accepted is not None and \
                   (highest_accepted is None or accepted.number > highest_accepted.number):
                    highest_accepted = accepted
            
            if len(promises) >= self.quorum_size:
                break
        
        if len(promises) < self.quorum_size:
            return None  # 无法获得多数派承诺
        
        # 决定提议值
        # 如果有已接受的值,必须使用它(保证一致性)
        proposal_value = highest_accepted.value if highest_accepted else value
        proposal = Proposal(proposal_number, proposal_value)
        
        # Phase 2: Accept
        acceptances = []
        
        for acceptor in promises:  # 只向承诺的 acceptor 发送
            if acceptor.accept(proposal):
                acceptances.append(acceptor)
            
            if len(acceptances) >= self.quorum_size:
                break
        
        if len(acceptances) >= self.quorum_size:
            return proposal_value  # 达成共识
        
        return None  # 未能获得多数派接受

# 使用示例
def demo_paxos():
    # 创建 5 个 Acceptor
    acceptors = [Acceptor(f"A{i}") for i in range(5)]
    quorum_size = 3  # 多数派
    
    # 创建 2 个 Proposer
    proposer1 = Proposer("P1", acceptors, quorum_size)
    proposer2 = Proposer("P2", acceptors, quorum_size)
    
    # Proposer 1 提议值 "X"
    result1 = proposer1.propose("X")
    print(f"Proposer 1 结果: {result1}")
    
    # Proposer 2 提议值 "Y"(应该被覆盖为 "X")
    result2 = proposer2.propose("Y")
    print(f"Proposer 2 结果: {result2}")
    
    # 验证所有 Acceptor 达成一致
    for acceptor in acceptors:
        accepted = acceptor.get_accepted()
        if accepted:
            print(f"{acceptor.acceptor_id} 接受了: {accepted.value}")

# demo_paxos()

2.5 Multi-Paxos:连续共识

基本 Paxos 只能对一个值达成共识。实际系统需要连续达成共识(如日志复制):

python
class MultiPaxos:
    """
    Multi-Paxos 实现(简化版)
    通过选主优化连续共识的性能
    """
    
    def __init__(self, node_id: str, peers: List[str]):
        self.node_id = node_id
        self.peers = peers
        self.log: Dict[int, any] = {}  # 日志:index -> value
        self.commit_index = 0
        self.current_leader: Optional[str] = None
        self.is_leader = False
    
    def run_leader_election(self) -> bool:
        """运行领导者选举(简化版)"""
        # 实际实现需要完整的 Paxos 过程
        # 这里简化为随机选择
        if random.random() < 0.5:
            self.is_leader = True
            self.current_leader = self.node_id
            return True
        return False
    
    def propose_log_entry(self, index: int, value: any) -> bool:
        """提议日志条目"""
        if not self.is_leader:
            # 转发给领导者
            return False
        
        # 使用 Paxos 达成共识
        # 实际实现需要与所有节点通信
        self.log[index] = value
        return True
    
    def append_entries(self, entries: List[tuple]):
        """领导者复制日志到跟随者"""
        for index, value in entries:
            self.log[index] = value

三、Raft算法:可理解的共识

3.1 为什么需要Raft?

Paxos 以其晦涩难懂著称。Lamport 的希腊寓言虽然生动,但实际实现 Paxos 极其困难。2014年,Diego Ongaro 和 John Ousterhout 提出了 Raft,目标是:

"提供与 Paxos 同等的安全性,同时比 Paxos 更容易理解和实现。"

3.2 领导者选举

Raft 通过强领导者简化共识过程:

Raft 角色:

1. Leader(领导者)
   └── 处理所有客户端请求
   └── 复制日志到 Follower
   └── 只有一个 Leader

2. Follower(跟随者)
   └── 被动接收日志
   └── 响应 Leader 请求
   └── 超时后转为 Candidate

3. Candidate(候选人)
   └── 发起领导者选举
   └── 获得多数票后成为 Leader

选举过程:

python
class RaftNode:
    """
    Raft 节点实现(核心逻辑)
    """
    
    def __init__(self, node_id: str, peers: List[str]):
        self.node_id = node_id
        self.peers = peers
        
        # 持久化状态
        self.current_term = 0
        self.voted_for = None
        self.log: List[LogEntry] = []
        
        # 易失状态
        self.state = "follower"  # follower, candidate, leader
        self.commit_index = 0
        self.last_applied = 0
        
        # Leader 状态
        self.next_index = {}
        self.match_index = {}
        
        # 定时器
        self.election_timeout = random.uniform(0.15, 0.3)  # 150-300ms
        self.heartbeat_interval = 0.05  # 50ms
    
    def start_election(self):
        """开始领导者选举"""
        self.state = "candidate"
        self.current_term += 1
        self.voted_for = self.node_id
        
        votes = 1  # 自己投自己
        
        # 向所有节点请求投票
        for peer in self.peers:
            if self.request_vote(peer):
                votes += 1
        
        # 获得多数票成为 Leader
        if votes > len(self.peers) / 2:
            self.become_leader()
    
    def request_vote(self, peer: str) -> bool:
        """
        向指定节点请求投票
        
        对方投票条件:
        1. term >= 自己的 current_term
        2. 日志至少和自己一样新
        """
        # 实际实现需要 RPC 调用
        # 这里简化处理
        return True
    
    def become_leader(self):
        """成为 Leader"""
        self.state = "leader"
        
        # 初始化 Leader 状态
        for peer in self.peers:
            self.next_index[peer] = len(self.log)
            self.match_index[peer] = 0
        
        # 立即发送心跳
        self.send_heartbeats()
    
    def send_heartbeats(self):
        """发送心跳(空 AppendEntries)"""
        for peer in self.peers:
            self.append_entries(peer, entries=[])

3.3 日志复制

Raft 的核心是日志复制:

python
@dataclass
class LogEntry:
    term: int
    index: int
    command: any

class RaftLogReplication:
    """
    Raft 日志复制实现
    """
    
    def __init__(self, node: RaftNode):
        self.node = node
    
    def append_entries(self, peer: str, entries: List[LogEntry]) -> bool:
        """
        Leader 向 Follower 复制日志
        
        参数:
        - prev_log_index: 前一个日志索引
        - prev_log_term: 前一个日志任期
        - entries: 要复制的日志条目
        - leader_commit: Leader 的 commit_index
        """
        prev_log_index = self.node.next_index[peer] - 1
        prev_log_term = 0
        
        if prev_log_index >= 0 and prev_log_index < len(self.node.log):
            prev_log_term = self.node.log[prev_log_index].term
        
        # 实际 RPC 调用
        # success = rpc_append_entries(peer, ...)
        
        # 处理响应
        # if success:
        #     self.node.match_index[peer] = prev_log_index + len(entries)
        #     self.node.next_index[peer] = self.node.match_index[peer] + 1
        # else:
        #     self.node.next_index[peer] -= 1  # 回退
        
        return True
    
    def handle_append_entries(self, term: int, leader_id: str,
                             prev_log_index: int, prev_log_term: int,
                             entries: List[LogEntry], leader_commit: int) -> bool:
        """
        Follower 处理 AppendEntries 请求
        
        返回 True 如果成功追加
        """
        # 1. 如果 term < current_term,拒绝
        if term < self.node.current_term:
            return False
        
        # 2. 重置选举定时器
        self.node.reset_election_timer()
        
        # 3. 如果 prev_log 不匹配,拒绝
        if prev_log_index >= 0:
            if prev_log_index >= len(self.node.log):
                return False
            if self.node.log[prev_log_index].term != prev_log_term:
                return False
        
        # 4. 追加新条目
        for i, entry in enumerate(entries):
            index = prev_log_index + 1 + i
            if index < len(self.node.log):
                # 冲突:删除现有条目及之后所有条目
                if self.node.log[index].term != entry.term:
                    self.node.log = self.node.log[:index]
                    self.node.log.append(entry)
            else:
                self.node.log.append(entry)
        
        # 5. 更新 commit_index
        if leader_commit > self.node.commit_index:
            self.node.commit_index = min(leader_commit, len(self.node.log) - 1)
        
        return True

3.4 安全性保证

Raft 通过以下机制保证安全性:

python
class RaftSafety:
    """
    Raft 安全性实现
    """
    
    def __init__(self, node: RaftNode):
        self.node = node
    
    def is_log_up_to_date(self, last_log_index: int, last_log_term: int) -> bool:
        """
        检查候选人的日志是否至少和自己一样新
        
        比较规则:
        1. 先比较最后条目的 term
        2. term 相同则比较 index
        """
        my_last_index = len(self.node.log) - 1
        my_last_term = 0
        
        if my_last_index >= 0:
            my_last_term = self.node.log[my_last_index].term
        
        if last_log_term != my_last_term:
            return last_log_term > my_last_term
        
        return last_log_index >= my_last_index
    
    def can_commit(self, index: int) -> bool:
        """
        检查日志条目是否可以提交
        
        条件:
        - 条目存储在多数派节点上
        - 条目的 term == current_term
        """
        if index >= len(self.node.log):
            return False
        
        entry = self.node.log[index]
        if entry.term != self.node.current_term:
            return False
        
        # 统计复制到多少节点
        replicated = 1  # Leader 自己
        for peer in self.node.peers:
            if self.node.match_index.get(peer, 0) >= index:
                replicated += 1
        
        return replicated > len(self.node.peers) / 2

四、安全性与活性证明

4.1 安全性(Safety)

选举安全性:任意任期内最多只有一个 Leader

证明(反证法):

假设任期 T 有两个 Leader:L1 和 L2

1. L1 成为 Leader 需要获得多数派投票
2. L2 成为 Leader 也需要获得多数派投票
3. 两个多数派必然有交集(鸽巢原理)
4. 交集中的节点在一个任期内只能投一次票
5. 矛盾!因此不可能有两个 Leader

日志匹配性:如果两个日志条目有相同的 index 和 term,则它们存储相同的命令

证明(归纳法):

基础:空日志满足条件

归纳:假设前 k 个条目满足条件
  考虑第 k+1 个条目
  - Leader 在一个 term 内只创建一个条目
  - 条目通过 AppendEntries 复制到 Follower
  - Follower 检查 prev_log 匹配后才追加
  - 因此所有节点的第 k+1 个条目相同

领导者完备性:如果日志条目已提交,则该条目会出现在所有未来 Leader 的日志中

证明:

1. 条目 E 已提交 → 存在于多数派节点
2. 新 Leader L 必须获得多数派投票
3. 投票条件:候选人的日志至少和自己一样新
4. 因此 L 的日志必然包含 E
   (否则投票给 L 的节点中至少有一个包含 E,
    该节点的日志比 L 新,不会投票给 L)

4.2 活性(Liveness)

选举活性: eventually 会选出一个 Leader

证明:

1. 所有节点的选举定时器随机化(150-300ms)
2. 最早超时的节点成为 Candidate
3. 如果没有冲突,该节点获得多数票成为 Leader
4. 如果有冲突(多个 Candidate),term 增加,重试
5. 由于超时随机, eventually 只有一个 Candidate 超时

日志复制活性: eventually 所有日志条目会被复制到所有节点

证明:

1. Leader 持续发送 AppendEntries(心跳)
2. 如果 Follower 日志不匹配,Leader 回退 next_index
3. eventually 找到匹配的 prev_log
4. 成功复制后更新 match_index
5. 重复直到所有节点同步

五、工程实践与优化

5.1 性能优化

批处理(Batching):

python
class BatchedRaft:
    """
    批处理优化
    """
    
    def __init__(self):
        self.pending_entries = []
        self.batch_timeout = 0.001  # 1ms
    
    def propose(self, command):
        """批量收集请求"""
        self.pending_entries.append(command)
        
        if len(self.pending_entries) >= 100:
            self.flush()
    
    def flush(self):
        """批量发送"""
        if not self.pending_entries:
            return
        
        # 一次性复制所有待处理条目
        self.replicate_batch(self.pending_entries)
        self.pending_entries = []

流水线(Pipelining):

python
class PipelinedRaft:
    """
    流水线优化:允许并发的 AppendEntries
    """
    
    def __init__(self):
        self.inflight_requests = 0
        self.max_inflight = 10
    
    def replicate(self, entry):
        """非阻塞复制"""
        if self.inflight_requests < self.max_inflight:
            self.inflight_requests += 1
            self.send_async(entry, self.on_response)
    
    def on_response(self, success):
        """异步回调"""
        self.inflight_requests -= 1
        # 处理响应...

5.2 成员变更

联合共识(Joint Consensus):

python
class MembershipChange:
    """
    Raft 成员变更(两阶段)
    """
    
    def __init__(self, raft: RaftNode):
        self.raft = raft
        self.old_config = set()
        self.new_config = set()
        self.joint_config = False
    
    def propose_membership_change(self, new_peers: List[str]):
        """提议成员变更"""
        # 阶段 1:切换到联合配置
        self.old_config = set(self.raft.peers)
        self.new_config = set(new_peers)
        self.joint_config = True
        
        # 在联合配置下达成共识
        self.raft.peers = list(self.old_config | self.new_config)
        
        # 阶段 2:切换到新配置
        if self.joint_committed():
            self.joint_config = False
            self.raft.peers = new_peers
    
    def joint_committed(self) -> bool:
        """检查联合配置是否已提交"""
        # 需要在旧配置和新配置都获得多数派
        pass

5.3 实际系统对比

系统 算法 特点 适用场景
etcd Raft 简单易用,生态丰富 配置存储、服务发现
Consul Raft 多数据中心支持 服务网格
TiKV Raft 分片(Region)+ Multi-Raft 分布式数据库
ZooKeeper ZAB 类似 Paxos 协调服务
Chubby Paxos Google 内部使用 分布式锁

5.4 共识算法选择框架

需求分析:

需要拜占庭容错?
├── 是 → PBFT / Tendermint / HotStuff
│   └── 区块链、加密货币
│
└── 否 → 节点数量?
    ├── 小规模(< 10)→ Raft
    │   └── 简单、易实现
    │
    └── 大规模(> 100)→ Multi-Paxos / EPaxos
        └── 高性能、低延迟

结语

从 Paxos 到 Raft,共识算法走过了三十年的发展历程。Paxos 奠定了理论基础,证明了分布式共识的可能性;Raft 则让这一理论变得可理解和可实现。

核心洞见:

  1. 多数派(Quorum) 是共识的核心机制,保证了安全性和容错性
  2. 领导者(Leader) 简化了共识过程,但引入了单点瓶颈
  3. 日志复制 是状态机复制的基础,保证了所有节点的状态一致性
  4. 安全性与活性的权衡 是分布式系统设计的永恒主题

理解共识算法,不仅是掌握一种技术,更是理解分布式系统的本质:在不确定性中寻找确定性,在混乱中建立秩序。


参考资源

经典论文:

  1. Lamport, L. (1998). "The Part-Time Parliament". ACM TOCS.
  2. Lamport, L. (2001). "Paxos Made Simple". ACM SIGACT News.
  3. Ongaro, D., & Ousterhout, J. (2014). "In Search of an Understandable Consensus Algorithm". USENIX ATC.
  4. Fischer, M. J., Lynch, N. A., & Paterson, M. S. (1985). "Impossibility of Distributed Consensus with One Faulty Process". JACM.

工程实践: 5. etcd Raft 实现:https://github.com/etcd-io/raft 6. Raft 可视化:http://thesecretlivesofdata.com/raft/ 7. Paxos 实现指南:https://martinfowler.com/articles/patterns-of-distributed-systems/

进阶阅读: 8. Howard, H., Malkhi, D., & Schwarzkopf, M. (2014). "Flexible Paxos: Quorum Intersection Revisited". 9. Moraru, I., Andersen, D. G., & Kaminsky, M. (2013). "There is More Consensus in Egalitarian Parliaments". SOSP.


创建时间:2026年04月11日
更新时间:2026年04月11日