共识算法:从Paxos到Raft的形式化理解与工程实践
一、问题定义:分布式共识的本质
1.1 什么是共识?
在分布式系统中,共识(Consensus) 是指多个节点就某个值达成一致的过程。看似简单的目标,在异步网络、节点故障、消息丢失的环境下,却成为计算机科学中最具挑战性的问题之一。
形式化定义:
给定 个进程组成的集合 ,每个进程 提出一个值 ,共识算法需要满足:
- 终止性(Termination):每个正确进程最终决定是否接受某个值
- 一致性(Agreement):所有正确进程接受相同的值
- 有效性(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 结果并不意味着共识不可能实现,而是告诉我们:
- 必须使用随机化(打破确定性)
- 或引入超时机制(将异步转为同步)
- 或容忍活性失效(牺牲终止性保证安全性)
二、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实现
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 只能对一个值达成共识。实际系统需要连续达成共识(如日志复制):
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
选举过程:
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 的核心是日志复制:
@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 通过以下机制保证安全性:
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):
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):
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):
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 则让这一理论变得可理解和可实现。
核心洞见:
- 多数派(Quorum) 是共识的核心机制,保证了安全性和容错性
- 领导者(Leader) 简化了共识过程,但引入了单点瓶颈
- 日志复制 是状态机复制的基础,保证了所有节点的状态一致性
- 安全性与活性的权衡 是分布式系统设计的永恒主题
理解共识算法,不仅是掌握一种技术,更是理解分布式系统的本质:在不确定性中寻找确定性,在混乱中建立秩序。
参考资源
经典论文:
- Lamport, L. (1998). "The Part-Time Parliament". ACM TOCS.
- Lamport, L. (2001). "Paxos Made Simple". ACM SIGACT News.
- Ongaro, D., & Ousterhout, J. (2014). "In Search of an Understandable Consensus Algorithm". USENIX ATC.
- 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日