跳到主要內容

系統設計

把資料結構放大到好幾台機器

主題 · Raft 共識與領導者選舉

Raft 共識與領導者選舉

好幾台機器要對同一串操作達成共識,即使其中幾台當機:Raft 先選出一個領導者,每筆寫入都要由它複製到多數節點才算數;領導者掛了就重新選一個。

情境
1/38
在 151 ms:
AT0BT1CT0DT0ET0
★ 領導者候選人跟隨者✕ 當機
每個節點的日誌
格1234
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 的資料庫,把共識留給它們。

和其他主題的關係

延伸閱讀
CAP 定理

出現在這些架構裡

時間與空間複雜度(Big O)

操作平均最差
一次寫入到提交
N 是節點數:送給每個跟隨者一次,等多數回覆;延遲約一個來回
O(N)O(N)
一輪選舉
要票訊息;分票時要重來好幾輪
O(N)O(N)
落後的跟隨者追上
L 是它缺的條目數:每次退一格重試(實作常一次退一整個任期)
O(L)O(L)

空間:O(L),每個節點存整份日誌;實際系統會定期做快照、截掉舊日誌

和其他做法比

已提交接受了卻沒提交被拒絕選舉沒有多數派領導者最久的提交等待
一切正常13 / 13001167 ms25 ms
領導者當機12 / 13012351 ms24 ms
領導者被隔離7 / 13602354 ms26 ms
兩個跟隨者當機13 / 13001167 ms27 ms
三個跟隨者當機13 / 130011767 ms1579 ms

每個情境模擬 4 秒、同一組種子;故障發生在 1.2 秒。「沒有多數派領導者」包含開機時第一次選舉的時間。所有情境、以及測試裡數百個隨機故障腳本,同一任期出現兩個領導者的次數都是 0,提交過的條目也從沒被改寫過。

領導者當機後恢復平均選舉次數
150–300 ms177 ms1.2
300–600 ms347 ms1.1
150–200 ms258 ms2.8
150–165 ms501 ms9.6
150–155 ms1560 ms13.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 個,卻多一個要等的對象。跨區部署時,每次提交至少要等到第二近的區域。
  • 任期、投票和日誌在回覆前必須寫進硬碟;重開機忘了自己投過票,同一任期就可能出現兩個領導者。