跳过正文

Raft Cluster Pre-Vote Protocol

作者
杨全烨
系统软件:操作系统、网络与分布式系统。
目录

本文章写于cluster bus 可插拔raft协议改造的早期,此时还未成形, 只是笔者学习的一个记录而已,因为笔者在尝试在 redis 的 cluster V2中实现Raft的预选举协议。 已经merge,我们还可以进行一些复盘。

原issue界面 这是valkey尝试实现可插拔的Raft协议中的预选举问题,这也算是一个有趣的工程实现。 这里我们还是需要首先研究:Diego Ongaro博士毕业论文原文 中的9.6节摘录: Preventing disruptions when a server rejoins the cluster One downside of Raft’s leader election algorithm is that a server that has been partitioned from the cluster is likely to cause a disruption when it regains connectivity. When a server is partitioned, it will not receive heartbeats. It will soon increment its term to start an election, although it won’t be able to collect enough votes to become leader. When the server regains connectivity sometime later, its larger term number will propagate to the rest of the cluster (either through the server’s RequestVote requests or through its AppendEntries response). This will force the cluster leader to step down, and a new election will have to take place to select a new leader. Fortunately, such events are likely to be rare, and each will only cause one leader to step down. If desired, Raft’s basic leader election algorithm can be extended with an additional phase to prevent such disruptions, forming the Pre-Vote algorithm. In the Pre-Vote algorithm, a candidate only increments its term if it first learns from a majority of the cluster that they would be willing to grant the candidate their votes (if the candidate’s log is sufficiently up-to-date, and the voters have not received heartbeats from a valid leader for at least a baseline election timeout). This was inspired by ZooKeeper’s algorithm, in which a server must receive a majority of votes before it calculates a new epoch and sends NewEpoch messages (however, in ZooKeeper servers do not solicit votes, other servers offer them). The Pre-Vote algorithm solves the issue of a partitioned server disrupting the cluster when it rejoins. While a server is partitioned, it won’t be able to increment its term, since it can’t receive permission from a majority of the cluster. Then, when it rejoins the cluster, it still won’t be able to increment its term, since the other servers will have been receiving regular heartbeats from the leader. Once the server receives a heartbeat from the leader itself, it will return to the follower state (in the same term). We recommend the Pre-Vote extension in deployments that would benefit from additional robustness. We also tested it in various leader election scenarios in AvailSim, and it does not appear to significantly harm election performance. 以下是翻译

问题背景:防止服务器重新加入集群时引起中断
#

Raft 领导者选举算法的一个缺点是,当一台因网络分区(Partition)而脱离集群的服务器恢复连接时,很可能会引起集群中断

当服务器被网络分区隔离时,它将无法接收到心跳(Heartbeats)。 它很快就会增加自己的任期号(Term)以发起一次选举(收不到来自于leader的心跳,所以自己会不断++term,然后超时,然后重新选举),尽管它不可能收集到足够的选票来成为领导者(Leader)。

一段时间后,当该服务器恢复网络连接时,它那更大的任期号将会传播到集群中的其他节点(通过该服务器发送的 RequestVote 请求,或者通过它对 AppendEntries 的响应)。

这里在这个断连的节点重新加入集群之后可以考察两种情况: 1.自己向其余的follower发送RequestVote请求,其余follower看到任期号之后提升自己的任期号,不会再回复当前的leader. 2.给leader回复的AppendEntries同理,leader会自动降级。

这将迫使集群当前的领导者退位,并且必须进行一次新的选举来选出新的领导者

你可能会认为,这个领导者就算term很大,但是日志不是很老么,能有什么影响? **但是实际的代码逻辑就是先比较term,如果对方的term更大,那么就提升自己的term,然后再看日志的新旧,造成的结果就是我的term提升了,但是日志上我拒绝了你,那么就是会造成大范围的term inflation,那么老领导者会被迫退位,造成新一轮的领导者选举,这一轮的term都会大幅上升。

逻辑时钟是我们的第一标准,难道你就不好奇,为什么我们不能在RequestVote的时候先检查日志的新旧?

实际上有这两点,原论文中有这样的规定:

论文 Figure 2 的全局规则(比 RequestVote 局部规则优先级更高)
#

这条规则不区分 RequestVote / AppendEntries,也不区分最后有没有 Grant 票。 这个是铁律,但是为什么是铁律?

若拒票时不升 term,会破坏「旧 leader 必须作废」的语义
#

假设允许:log 不够新 → 直接 return,term 不变。

Follower F:term = 5
过期 Leader L:term = 5(多数派其实已失联,但 L 自己不知道)
捣乱者 C:  term = 6,log 很旧

F 收到 C 的 RequestVote:log 旧 → 拒票,term 仍为 5。
L 继续用 term 5 发心跳,F 仍接受。

从可用性看似乎更好,但和论文设计目标冲突:

  • 论文要求:一旦集群里出现 term 6 的选举信号,所有节点都应认为「epoch 6 已开始」,term 5 的 leader 不能再被当作合法 leader。
  • 安全证明依赖 term 作为全集群单调的逻辑时钟:见过 term T 的节点,不会再把 < T 的 leader 当真。

若部分节点升到 6、部分因「log-first 拒票」留在 5, 会出现 term 视图分裂,AppendEntries / RequestVote 行为不一致,证明假设被破坏

很简单,全局的逻辑时钟必须作为第一准则

原生的go代码大概是这样的:

// 我是follower,当有candidate想让我投票的时候,我该如何处理?
// 这个视角应该就是我看args,然后填reply即可
// 在相同任期下,你才能进行投票的操作
func (rf *Raft) RequestVote(args *RequestVoteArgs, reply *RequestVoteReply) {
	rf.mu.Lock()
	defer rf.mu.Unlock()

	// 任期比我更小,不能给你投票
	if args.Term < rf.currentTerm {
		reply.Term = rf.currentTerm
		reply.VoteGranted = false
		return
	}

	// 任期比我更大,我先升级一下自己
	if args.Term > rf.currentTerm {
		rf.currentState = follower
		rf.currentTerm = args.Term
		rf.votedFor = -1 // 清理一下本次term的投票

		rf.persist()
	}

	// 任期和我相同(包括刚升级结束的)
	// 我投票了,但是本来就没给你投
	if rf.votedFor != -1 && rf.votedFor != args.CandidateId {
		reply.Term = rf.currentTerm
		reply.VoteGranted = false
		return
	}

	// 处理log,是否够新,这里我们先看Term,然后看Log
	// 实际工程有预投票的实现,比较有意思,我们这里暂时不管
	if !rf.isLogUpToDate(args.LastLogIndex, int(args.LastLogTerm)) {
		reply.Term = rf.currentTerm
		reply.VoteGranted = false
		return
	}

	// ok
	rf.votedFor = args.CandidateId

	// 因为进行了一次投票所以需要持久化
	rf.persist()
	// 这次选票成功,也认为是一次压制
	rf.lastRPC = time.Now()
	reply.VoteGranted = true
	reply.Term = rf.currentTerm
}

幸运的是,此类事件通常很少发生,并且每次发生只会导致一位领导者退位。 如果需要,Raft 的基础领导者选举算法可以通过增加一个额外的阶段来进行扩展,以防止这种中断,从而形成了预投票(Pre-Vote)算法.

在预投票算法中,候选人(Candidate)只有在首先从集群的大多数节点那里得知它们愿意将选票投给自己时,才会增加自己的任期号(前提是该候选人的日志足够新,并且投票者在至少一个基准选举超时时间内,没有收到来自有效领导者的心跳,就是leader应该差不多还活着呢,你为啥要让我投票)。

这个基准超时时间不能随机化,就定成一个固定的时间最好。

这受到了 ZooKeeper 算法的启发,在 ZooKeeper 中,服务器必须在计算新纪元(Epoch)并发送 NewEpoch 消息之前获得大多数的选票(不过,在 ZooKeeper 中服务器并不主动拉票,而是由其他服务器主动提供选票)。

预投票算法解决了被分区的服务器在重新加入时破坏集群稳定性的问题。 当服务器处于分区状态时,由于无法获得集群大多数节点的许可,它将无法增加自己的任期号

注意,这里是直接无法增加,而不是增加了但是无法影响。

随后,当它重新加入集群时,它仍然无法增加自己的任期号,因为其他服务器一直在正常接收来自当前领导者的心跳一旦该服务器接收到来自当前领导者自身的心跳,它就会恢复到跟随者(Follower)状态(并且保持在相同的任期内)。

对于那些要求更高鲁棒性(Robustness)的部署环境,我们推荐使用预投票扩展机制

我们还在 AvailSim 的各种领导者选举场景中对其进行了测试,结果表明它似乎不会对选举性能造成明显的负面影响。(所以这个究竟是否需要引入还是需要评估性能问题,就是会造成一点负面影响,但是影响不大)

当我们理解了问题了之后,就该看看代码究竟是如何组织的? 不管是什么系统使用了Raft算法,我们都需要注意以下几个点。

在这个仓库内部: 1.RPC是如何进行调用的? 2.日志,server和Raft成员相关的结构体定义在哪里?角色状态是在哪里定义的? 3.诸如AppendEntriesRequestVote这样的RPC调用定义在哪里? 4.触发选举,心跳等方法在哪里? 5.随机超时时间是在哪里进行定义的? 等等

Raft架构分析
#

怎么通信?
#

首先因为是在Redis内部,我们不会使用现成的RPC框架。 cluster bus TCP 连接 + 自定义二进制头 + 文本消息 (我就说其实RPC也没什么特殊的) 这里的「RPC」= 在 clusterLink 上发/收一条 bus 消息,流程如下。

当发送出站消息的时候:

消息都在16379上面进行传输。

业务逻辑 → wireNewMsg("AE"|"VOTE_REQ"|...) → wireFinishMsg → clusterRaftSendMsg
         → clusterLinkSendBlock → TCP 写出

比如:

// 比如在总线上发送一个sds消息,如果发送失败,就直接free sds
static void clusterRaftSendMsg(clusterLink *link, sds msg) {
    if (!link || !link->conn) {
        sdsfree(msg);
        return;
    }
    size_t len = sdslen(msg);
    clusterMsgSendBlock *block = clusterAllocMsgSendBlock(len);
    memcpy(block->data, msg, len);
    sdsfree(msg);
    clusterLinkSendBlock(link, block);
    clusterMsgSendBlockDecrRefCount(block);
    RAFT_STATE()->stats_bytes_sent += len;
}

广播用 clusterRaftBroadcast / clusterRaftBroadcastAppendEntries, 遍历 server.cluster->nodes 对每个有 link 的 peer 发。

当接收入站消息的时候:

事件循环读 TCP → clusterReadHandler,首先读header, 然后检查内存,了解到整个数据包都进入内存了之后调用对应的方法来进行处理。

  1. 读满 RCVBUF_MIN_READ_LEN 字节
  2. clusterCurrentBus->validateMessageHeader() — Raft 下检查 "RAFT" 魔数 + 大端 totlen
  3. 读满整包后 clusterCurrentBus->processMessage(link)
  4. Raft 里即 clusterRaftProcessMessage,按首 token 分派到各 clusterRaftProcess* 处理函数

cluster_link.c:

        if (rcvbuflen >= RCVBUF_MIN_READ_LEN) {
            uint32_t totlen = clusterCurrentBus->validateMessageHeader(link->rcvbuf);
            if (totlen && rcvbuflen == totlen) {
	            // 分派给对应的方法进行处理
                if (clusterCurrentBus->processMessage(link)) {

还有一条路径就是: 当follower收到来自于client的某些消息的时候会进行相应的转发操作,转给leader处理。

比如像是写操作

结构体是如何定义的?
#

因为是日志复制状态机,我们的日志是如何定义的:

// 存在哪些entry类型?
enum raftEntryType {
	RAFT_ENTRY_NODE_JOIN = 1, /* Add a node 增加一个节点 */ 
    RAFT_ENTRY_NODE_FORGET = 2, /* Remove a node 移除一个节点 */
    RAFT_ENTRY_SLOT_CHANGE = 3, /* Slot ownership 占有一个slot */
    RAFT_ENTRY_SET_REPLICA_OF = 4, /* Replication topology 复制拓扑 */
    RAFT_ENTRY_FAILOVER = 5, /* Failover (manual or automatic) 故障转移 */
    RAFT_ENTRY_NODE_INFO = 6, /* IP, port, hostname, etc. 一些节点信息 */
    RAFT_ENTRY_NODE_FAIL = 7, /* Node failure detected 挂了 */
    RAFT_ENTRY_NODE_RECOVER = 8, /* Node recovery detected 节点恢复 */
    ......
};


typedef struct {
    uint64_t term;  // 任期
    uint64_t index; // 日志索引
    uint8_t type; /* enum raftEntryType 当前条目的基本类型 */
    sds data; /* Space-separated command arguments 空格隔开的传参 */ 
}
raftLogEntry;

对于一个Raft Node存在的多种角色:

enum raftRole {
    RAFT_ROLE_JOINER = 0, /* Stepped down singleton waiting to be added to a cluster */
    RAFT_ROLE_FOLLOWER = 1,
    RAFT_ROLE_CANDIDATE = 2,
    RAFT_ROLE_LEADER = 3,
};

这个JOINER就是单机等待加入cluster的一个中间态

还有就是全局状态的定义,比如像是一些持久化的状态啊,leader需要维护的状态之类的: (这里所谓的全局应该就是leader的状态?)

/* --------------------------------------------------------------------------
 * Protocol-specific state (stored in clusterState.protocol_data)
 * -------------------------------------------------------------------------- */

typedef struct {
    /* Persistent Raft state */
    uint64_t current_term;
    char voted_for[CLUSTER_NAMELEN]; /* All zeros = none */

    /* Volatile Raft state */
    enum raftRole role;
    uint64_t commit_index;
    uint64_t last_applied;

    /* Log */
    raftLogEntry ** log;
    uint64_t log_count;
    uint64_t log_alloc;

    /* Pending proposals waiting for commit. */
    list * pending_proposals; /* list of raftPendingProposal */

    /* Pending MEET callbacks waiting for NODE_JOIN commit. */
    list * pending_meets; /* CLUSTER MEET commands waiting for OK reply */
    list * deferred_meets; /* Inbound MEET messages deferred until size > 1 */

    /* NODE_INFO divergence detection. */
    sds my_last_committed_info;
    mstime_t last_node_info_check;

    /* Election */
    int votes_received;
    mstime_t election_timeout; /* Randomized timeout */
    mstime_t last_heartbeat; /* Last time we heard from leader */
    mstime_t last_repl_offsets_broadcast; /* Last REPL_OFFSETS broadcast */
    mstime_t joiner_since; /* When we became a joiner (for timeout) */

    /* Leader identity */
    char leader[CLUSTER_NAMELEN]; /* All zeros = unknown */

    /* Manual failover state (on the replica side) */
    mstime_t mf_end; /* Timeout for manual failover, 0 = not in progress */
    void * mf_ctx; /* blockedAsyncHandle for the CLUSTER FAILOVER client */
    void( * mf_callback)(void * ctx,
        const char * error);

    /* Automatic failover state (on the replica side) */
    mstime_t failover_time; /* When to propose FAILOVER based on rank, 0 = inactive */

    /* Message stats for CLUSTER INFO */
    long long stats_module_messages_sent;
    long long stats_module_messages_received;
    long long stats_publish_messages_sent;
    long long stats_publish_messages_received;
    uint64_t stats_bytes_sent;
    uint64_t stats_bytes_received;
    uint64_t stats_pubsub_bytes_sent;
    uint64_t stats_pubsub_bytes_received;
    uint64_t stats_module_bytes_sent;
    uint64_t stats_module_bytes_received;
}
clusterRaftState;

差不多也能看懂一些经典的成员了。

接下来就是每个Raft Node 所持有的状态:

typedef struct {
    uint64_t next_index;
    uint64_t match_index;   // 记录日志的复制进度 论文5.3
    mstime_t last_ack_time; /* Last time we received AE_ACK from this peer 用于NODE_FAIL的检测,超过固定时间认为就会下线 */
    unsigned int pending_fail_change: 1; /* NODE_FAIL or NODE_RECOVER in flight */
}
clusterNodeRaftData;

在这里,那么一个valkey的进程就是一个server.

补充一下,AE就是AppendEntries的简称,不是事件

AppendEntries / RequestVote
#

上面这两个经典的RPC呢

没有独立 .proto 或函数指针表;消息类型是字符串常量, 在 clusterRaftProcessMessage 里 strcasecmp 分发(2026–2047 行):

论文 RPC 线协议名 发送 处理
AppendEntries AE / AE_ACK clusterRaftSendAppendEntries / clusterRaftSendAppendEntriesResponse clusterRaftProcessAppendEntries / ...Response
RequestVote VOTE_REQ / VOTE clusterRaftSendRequestVote / clusterRaftSendVoteResponse clusterRaftProcessRequestVote / ...Response

可以看看AE的消息格式:

 * Header line: AE <leader-id> <term> <prev-log-idx> <prev-log-term> <commit> <count>
 * Entry lines: <term> <type> <data>

相当于字符串传递,然后进行消息解析。 if count == 0 那本条消息就应该是一个心跳。

话不多说,接下来我们尝试开始实现预投票的机制。 你可以在issue界面看到我们的实现原理,此处便不赘述。

一些细节的工程考虑
#

当candidate选举超时的时候
#

在拿下了pre-vote的选举成功之后,此时这个机器发生了partitioned,如若一直保持着candidate的状态,此时还是会发生term inflation,当此机器回归的时候,还是会带着很高的term,同时身份还是candidate,出现term inflation,所以每次选举超时都需要重新发起pre-vote是合理的举措. 这意味着每次正式选举超时都需要重新发起一次预选举.

理解pre-vote机制的核心
#

pre vote的核心是引入中间状态pre-candidate,在这个状态下,current term根本就不会无限膨胀,就是不会被持久化.

stateDiagram-v2
    [*] --> Follower
    Follower --> PreCandidate: election timeout
    PreCandidate --> Candidate: majority PRE_VOTE grants
    PreCandidate --> Follower: pre-vote timeout (no quorum)
    Candidate --> Leader: majority VOTE grants
    Candidate --> Follower: election lost / higher term seen
    Leader --> Follower: step down (higher term)
    Follower --> Candidate: TIMEOUT_NOW (leader transfer)
    note right of PreCandidate
        term NOT incremented
    end note
    note right of Candidate
        term incremented
    end note

上面就是我们在PR页面做的mermaid表。 Pre-vote 的本质就是多一个 PRE_CANDIDATE 中间态,把「试探能不能选」和「真正抬 term 开正式选举」拆开:

阶段 current_term 是否持久化新 term
PRE_CANDIDATE(pre-vote) 不变 否(只是线上发 current_term+1 去试探)
StartElection() 之后 ++ 是(todo_save_config 会落盘)

所以在 partition 里反复超时、反复 pre-vote 失败时,节点会在 PRE_CANDIDATE ↔ FOLLOWER 之间转, term 不会因为“选不上”而一层层叠高——这正是 §9.6 要防的 disruption。

测试脚本
#

这里测试的学问很大,值得认真研究,测试的系统思想是我们所欠缺的,同时也是AI编程时代很重要的一点。

cd /home/ada/Project/valkey

# 确保二进制最新
make -C src valkey-server -j$(nproc)

# 单测(用不同 baseport 避免 21111 段冲突)
./runtest --single tests/unit/cluster/cluster-raft-prevote.tcl --baseport 21200
./runtest --single tests/unit/cluster/cluster-raft-proto.tcl --baseport 21300

1.复用性 2.有效性 3.覆盖率问题 如果你想保证自己的代码有效,首先需要了解如何进行测试

此类测试的整体架构:

flowchart TB
    subgraph CI["CI 两条流水线"]
        A["test-cluster-raft job<br/>./runtest --cluster-raft --single tests/unit/cluster"]
        B["codecov job<br/>make lcov → 全量 runtest + gtest"]
    end

    subgraph BlackBox["黑盒测试层 (Tcl)"]
        P["cluster-raft-proto.tcl<br/>伪造 peer + TCP cluster bus"]
        E["cluster-raft-prevote.tcl<br/>3 真实节点 + SIGSTOP"]
        R["cluster-raft.tcl<br/>quorum freshness"]
    end

    subgraph Server["valkey-server 进程"]
        EP["epoll 读 cluster bus fd"]
        PARSE["RAFT 帧解析 → argv[]"]
        DISPATCH["clusterRaftProcessMsg()"]
        RAFT["cluster_raft.c pre-vote 状态机"]
    end

    A --> P & E & R
    B --> P & E & R
    P -->|"socket :port+10000"| EP
    E -->|"Redis 协议 CLUSTER INFO / pause_process"| RAFT
    EP --> PARSE --> DISPATCH --> RAFT

协议黑盒:Leader lease 活跃时拒绝 PRE_VOTE
#

sequenceDiagram
    participant Tcl as Tcl fake peer
    participant Bus as cluster bus :port+10000
    participant S as valkey-server (singleton leader)

    Tcl->>Bus: HELLO
    Bus->>Tcl: HI
    Note over S: leader=自己, last_heartbeat 新鲜<br/>lease 有效
    Tcl->>Bus: PRE_VOTE_REQ term+1 log 0/0
    Bus->>S: clusterRaftProcessPreVoteRequest()
    S->>S: clusterRaftCanGrantVote() → 0 (lease)
    Bus->>Tcl: PRE_VOTE current_term 0
    Note over S: current_term 不变

协议黑盒:无 lease 时 grant PRE_VOTE
#

sequenceDiagram
    participant Tcl1 as fake candidate1
    participant Tcl2 as fake candidate2
    participant S as valkey-server

    Tcl1->>S: VOTE_REQ term+1
    Note over S: step down → follower, leader=""
    Tcl2->>S: PRE_VOTE_REQ term+2 log 0/0
    S->>S: CanGrantVote() → 1 (无 leader + log OK)
    S->>Tcl2: PRE_VOTE term+2 1
    Note over S: current_term 仍为 vote_term,未 inflate

协议黑盒:候选者 log 过旧时拒绝
#

flowchart TD
    A["AE 写入 NODE_JOIN<br/>receiver last_log ≥ 1"] --> B["VOTE_REQ step-down<br/>清 leader lease"]
    B --> C["PRE_VOTE_REQ log 0/0<br/>(stale)"]
    C --> D["CanGrantVote: candidate_last_term < my_last_term"]
    D --> E["PRE_VOTE deny, term 不变"]

协议黑盒:pre-vote 超时、无 quorum 不 inflate term(最关键)
#

sequenceDiagram
    participant FakeL as Tcl fake leader (inbound)
    participant S as valkey-server (follower)
    participant FakeP as Tcl listen (outbound link)

    FakeL->>S: AE term=2, NODE_JOIN x2
    Note over S: role=follower, term=2, cluster_size=2
    S->>FakeP: HELLO (outbound node->link)
    FakeP->>S: HI
    Note over FakeL: 停止 AE → last_heartbeat 过期
    S->>S: clusterRaftCron → clusterRaftStartPreVote()
    S->>FakeP: PRE_VOTE_REQ term=3
    Note over S: term 仍为 2 ✓
    Note over FakeP: 故意不回复 (无 quorum)
    S->>FakeP: PRE_VOTE_REQ term=3 (第二轮)
    FakeP->>S: PRE_VOTE 3 1 (grant)
    S->>S: pre_votes_received++ → quorum → StartElection()
    S->>FakeP: VOTE_REQ term=3
    Note over S: term=3 ✓

端到端黑盒:Leader 故障 → pre-vote → 选举
#

这就是一个正常选举的逻辑。

sequenceDiagram
    participant N0 as Node0 (leader)
    participant N1 as Node1
    participant N2 as Node2

    Note over N0,N2: start_cluster 3, term=T
    N0->>N0: pause_process (SIGSTOP)
    Note over N0: 进程冻结,不处理 AE<br/>last_heartbeat 在 N1/N2 过期
    N1->>N2: PRE_VOTE_REQ (经真实 cluster bus)
    N2->>N1: PRE_VOTE grant
    N1->>N1: StartElection, term=T+1
    N1->>N2: VOTE_REQ ...
    Note over N1: 日志含 "Starting Raft pre-vote"

其实还是测试一下关键的路径即可。