Raft 共識與領導者選舉
好幾台機器要對同一串操作達成共識,即使其中幾台當機:Raft 先選出一個領導者,每筆寫入都要由它複製到多數節點才算數;領導者掛了就重新選一個。
情境
1/38
在 151 ms:
★ 領導者候選人跟隨者✕ 當機
| 格 | 1 | 2 | 3 | 4 |
|---|---|---|---|---|
| A | ||||
| B | ||||
| C | ||||
| D | ||||
| E |
已提交(就這個節點所知)尚未提交T:寫入時的任期
150 ms:B 的計時器到期前沒收到領導者的訊息。它進入第 1 任期、投自己一票,再向其他四個節點要票。
已提交的寫入
12 / 13
接受了卻沒提交
0
選舉次數
2
同一任期兩個領導者
0
模型:訊息延遲 5–15 ms,跨過斷線或送到當機節點就遺失;選舉計時器在 150–300 ms 之間隨機;用戶端每 250 ms 寫一次,送到它最後知道的領導者。未模擬:磁碟延遲、快照、成員變更,以及用戶端重送沒提交的寫入。
亮起來的是這一步執行的程式碼
type Entry = { term: number; cmd: string };type Msg = { type: string; from: number; to: number; term: number; [field: string]: any }; class RaftNode { role = "follower"; term = 0; votedFor: number | null = null; log: Entry[] = []; commitIndex = 0; votes = new Set<number>(); nextIndex: number[]; matchIndex: number[]; constructor(public id: number, public n: number) { this.nextIndex = new Array(n).fill(1); this.matchIndex = new Array(n).fill(0); } lastIndex() { return this.log.length; } lastTerm() { return this.log.length ? this.log[this.log.length - 1].term : 0; } stepDown(term: number) { if (term <= this.term) return; this.term = term; this.role = "follower"; this.votedFor = null; } // The election timer ran out without hearing from a leader. startElection(): Msg[] { this.term += 1; this.role = "candidate"; this.votedFor = this.id; this.votes = new Set([this.id]); const out: Msg[] = []; for (let p = 0; p < this.n; p++) { if (p !== this.id) out.push({ type: "vote", from: this.id, to: p, term: this.term, lastIndex: this.lastIndex(), lastTerm: this.lastTerm() }); } return out; } onRequestVote(m: Msg): Msg { this.stepDown(m.term); const upToDate = m.lastTerm > this.lastTerm() || (m.lastTerm === this.lastTerm() && m.lastIndex >= this.lastIndex()); const granted = m.term === this.term && (this.votedFor === null || this.votedFor === m.from) && upToDate; if (granted) this.votedFor = m.from; return { type: "voteReply", from: this.id, to: m.from, term: this.term, granted }; } onVoteReply(m: Msg): Msg[] { this.stepDown(m.term); if (this.role !== "candidate" || m.term !== this.term || !m.granted) return []; this.votes.add(m.from); if (this.votes.size * 2 <= this.n) return []; this.role = "leader"; this.nextIndex.fill(this.lastIndex() + 1); this.matchIndex.fill(0); this.matchIndex[this.id] = this.lastIndex(); return this.heartbeat(); } propose(cmd: string): Msg[] | null { if (this.role !== "leader") return null; // only a leader takes writes this.log.push({ term: this.term, cmd }); this.matchIndex[this.id] = this.lastIndex(); return this.heartbeat(); } heartbeat(): Msg[] { const out: Msg[] = []; for (let p = 0; p < this.n; p++) if (p !== this.id) out.push(this.appendFor(p)); return out; } appendFor(p: number): Msg { const prevIndex = this.nextIndex[p] - 1; const prevTerm = prevIndex > 0 ? this.log[prevIndex - 1].term : 0; return { type: "append", from: this.id, to: p, term: this.term, prevIndex, prevTerm, entries: this.log.slice(prevIndex), leaderCommit: this.commitIndex }; } onAppendEntries(m: Msg): Msg { this.stepDown(m.term); const reply = (success: boolean, matchIndex: number) => ({ type: "appendReply", from: this.id, to: m.from, term: this.term, success, matchIndex }); if (m.term < this.term) return reply(false, 0); // a deposed leader this.role = "follower"; if (m.prevIndex > this.lastIndex() || (m.prevIndex > 0 && this.log[m.prevIndex - 1].term !== m.prevTerm)) { return reply(false, 0); } for (let k = 0; k < m.entries.length; k++) { const index = m.prevIndex + 1 + k; if (index <= this.lastIndex() && this.log[index - 1].term !== m.entries[k].term) { this.log.length = index - 1; } if (index > this.lastIndex()) this.log.push(m.entries[k]); } const match = m.prevIndex + m.entries.length; if (m.leaderCommit > this.commitIndex) { this.commitIndex = Math.max(this.commitIndex, Math.min(m.leaderCommit, match)); } return reply(true, match); } onAppendReply(m: Msg): Msg[] { this.stepDown(m.term); if (this.role !== "leader" || m.term !== this.term) return []; if (!m.success) { this.nextIndex[m.from] = Math.max(1, this.nextIndex[m.from] - 1); // back up and retry return [this.appendFor(m.from)]; } this.matchIndex[m.from] = Math.max(this.matchIndex[m.from], m.matchIndex); this.nextIndex[m.from] = this.matchIndex[m.from] + 1; // Commit the newest entry of this term that a majority holds. for (let index = this.lastIndex(); index > this.commitIndex; index--) { if (this.log[index - 1].term !== this.term) break; const holders = this.matchIndex.filter((x) => x >= index).length; if (holders * 2 > this.n) { this.commitIndex = index; break; } } return []; } restart() { // term, vote and log were on disk this.role = "follower"; this.votes = new Set(); } handle(m: Msg): Msg[] { if (m.type === "vote") return [this.onRequestVote(m)]; if (m.type === "voteReply") return this.onVoteReply(m); if (m.type === "append") return [this.onAppendEntries(m)]; return this.onAppendReply(m); }}模型假設與範圍
- 這是可重現的教學模型;延遲、容量、故障率與工作負載是設定或樣本,不能直接當作正式系統的效能承諾。
- 五節點固定成員、5–15 ms 訊息延遲、隨機選舉逾時;當機保留日誌與任期。未模擬磁碟延遲、快照、成員變更與 Byzantine 故障。
小挑戰
五個節點中三個跟隨者當機了。選一個恢復動作,讓 1700 ms 之後的新寫入可以提交。
恢復動作
恢復後提交的新寫入
0
調整選項,再檢查結果。
提示
目前只有兩個節點存活;修復網路不會讓當機的節點醒來。
查看解答
重新啟動當機節點,恢復至少三個可互通節點。多數重新選出領導者後才可提交新寫入。
什麼時候用
- 幾個節點必須對同一串決定的順序完全一致的時候:誰是主節點、設定檔、分片對照表、鎖。Raft 把它拆成三件事:選出一個領導者、領導者把日誌複製出去、多數存好了才算提交。
- 要能容忍少數節點故障:5 個節點可以壞 2 個、3 個可以壞 1 個。壞掉超過一半就停下來不收寫入,而不是冒險分裂成兩個版本。
- 通常不用自己寫:用 etcd、ZooKeeper、Consul 這類協調服務,或內建 Raft 的資料庫,把共識留給它們。
和其他主題的關係
- 由這些組成
- 一致性與 Quorum資料庫複製
- 延伸閱讀
- CAP 定理
出現在這些架構裡
時間與空間複雜度(Big O)
| 操作 | 平均 | 最差 |
|---|---|---|
| 一次寫入到提交 N 是節點數:送給每個跟隨者一次,等多數回覆;延遲約一個來回 | O(N) | O(N) |
| 一輪選舉 要票訊息;分票時要重來好幾輪 | O(N) | O(N) |
| 落後的跟隨者追上 L 是它缺的條目數:每次退一格重試(實作常一次退一整個任期) | O(L) | O(L) |
空間:O(L),每個節點存整份日誌;實際系統會定期做快照、截掉舊日誌
和其他做法比
| 已提交 | 接受了卻沒提交 | 被拒絕 | 選舉 | 沒有多數派領導者 | 最久的提交等待 | |
|---|---|---|---|---|---|---|
| 一切正常 | 13 / 13 | 0 | 0 | 1 | 167 ms | 25 ms |
| 領導者當機 | 12 / 13 | 0 | 1 | 2 | 351 ms | 24 ms |
| 領導者被隔離 | 7 / 13 | 6 | 0 | 2 | 354 ms | 26 ms |
| 兩個跟隨者當機 | 13 / 13 | 0 | 0 | 1 | 167 ms | 27 ms |
| 三個跟隨者當機 | 13 / 13 | 0 | 0 | 1 | 1767 ms | 1579 ms |
每個情境模擬 4 秒、同一組種子;故障發生在 1.2 秒。「沒有多數派領導者」包含開機時第一次選舉的時間。所有情境、以及測試裡數百個隨機故障腳本,同一任期出現兩個領導者的次數都是 0,提交過的條目也從沒被改寫過。
| 領導者當機後恢復 | 平均選舉次數 | |
|---|---|---|
| 150–300 ms | 177 ms | 1.2 |
| 300–600 ms | 347 ms | 1.1 |
| 150–200 ms | 258 ms | 2.8 |
| 150–165 ms | 501 ms | 9.6 |
| 150–155 ms | 1560 ms | 13.0 |
選舉計時器的隨機範圍,各跑 12 個種子取平均:領導者在 1 秒時當機,量到新領導者出現為止。範圍太窄,好幾個節點會同時發起選舉、把票分掉,只好一輪一輪重來。
真實世界裡的它
- etcd(Kubernetes 的所有狀態都存在這裡)、Consul 都直接用 Raft。
- CockroachDB、TiKV 把資料切成很多段,每段各跑一組 Raft;Kafka 的 KRaft 模式用 Raft 取代了 ZooKeeper。
- ZooKeeper 的 ZAB(ZooKeeper Atomic Broadcast)和 Paxos 解決同一個問題;Raft 是 2014 年為了「好懂」而設計的。
取捨與陷阱
- 被隔離的舊領導者還會收寫入:它不知道自己被取代了,接下的寫入永遠湊不到多數(上面「領導者被隔離」有 6 筆)。用戶端必須等提交才算成功,不能只看領導者有沒有收下。
- 讀取也要小心:直接讀領導者的本地資料,可能讀到被隔離的舊領導者的過期值。要嘛把讀取也走一次日誌,要嘛先確認自己仍是多數派的領導者(ReadIndex、租約)。
- 只能提交自己任期的條目:新領導者不能因為舊任期的條目已經在多數節點上就直接提交它,要等自己任期的一筆條目提交時順便帶過去。這是論文圖 8 的陷阱。
- 節點數取奇數:4 個節點和 3 個一樣只能壞 1 個,卻多一個要等的對象。跨區部署時,每次提交至少要等到第二近的區域。
- 任期、投票和日誌在回覆前必須寫進硬碟;重開機忘了自己投過票,同一任期就可能出現兩個領導者。