Nacos 2.x 源码深度解析 (十五):JRaft 架构 ——CP 集群 Leader 选举机制

《Nacos 2.x源码深度解析》专栏目录
一、架构通信篇:
《Nacos 2.x 源码深度解析 (一):架构整体全貌 —— 核心模块划分与版本演进》
《Nacos 2.x 源码深度解析 (二):通信协议迭代 —— HTTP长轮询到gRPC演进》
二、配置中心篇
《Nacos 2.x 源码深度解析 (三):配置中心客户端 —— 启动加载与自动装配》
《Nacos 2.x 源码深度解析 (四):配置中心服务端 —— 事件总线与数据持久化》
《Nacos 2.x 源码深度解析 (五):gRPC 推送链路 —— 配置变更下发与动态刷新》
《Nacos 2.x 源码深度解析 (六):三级缓存体系 —— 降级兜底与故障自愈机制》
三、服务注册发现篇
《Nacos 2.x 源码深度解析 (七):服务注册流程 —— 客户端上报与服务端存储》
《Nacos 2.x 源码深度解析 (八):服务订阅机制 —— 从首次订阅到gRPC双向流变更通知》
四、grpc连接内核篇
《Nacos 2.x 源码深度解析 (九):双向流设计 —— 连接创建复用与销毁》
《Nacos 2.x 源码深度解析 (十):心跳保活策略 —— 断线检测与重连源码》
《Nacos 2.x 源码深度解析 (十一):RPC 请求调度 —— 收发模型与线程池处理》
五、集群一致性篇
《Nacos 2.x 源码深度解析 (十二):集群基础交互 —— 节点感知与基础数据同步》
《Nacos 2.x 源码深度解析 (十三):Distro 协议 ——AP 模式异步数据同步原理》
《Nacos 2.x 源码深度解析 (十四):Distro 容错处理 —— 数据校验与冲突修复》
《Nacos 2.x 源码深度解析 (十五):JRaft 架构 ——CP 集群 Leader 选举机制》
《Nacos 2.x 源码深度解析 (十六):JRaft 日志复制 —— 同步规则与快照管理》
文章目录
在上一篇文章中,我们深入分析了 Distro 协议的容错处理机制,理解了定时校验如何发现异步同步可能遗漏的数据差异,以及校验失败后如何通过补偿同步修复。
但在 Nacos 集群一致性体系中,持久化实例和配置数据的一致性保障究竟依靠什么?JRaft 协议作为 Nacos 的 CP 一致性实现,其整体架构如何设计与分层?Multi-Raft Group 机制为何要将数据按功能模块划分为独立的 Raft 组,每个组独立选举 Leader?Leader 选举过程中的心跳超时(electionTimeout)、onLeaderStart/onLeaderStop 回调以及 RaftEvent 元数据更新是如何形成闭环的?读写请求又如何根据节点角色路由到正确的 Leader?本文将切换视角,从 AP 转向 CP,以 JRaft 协议的整体架构为起点,深入 Multi-Raft Group 的创建、Leader 选举的内部流转,再到读写请求的路由策略,完整剖析 Nacos 中 CP 一致性的 Leader 选举与请求路由机制。
一、CP 协议在 Nacos 中的定位
在深入 JRaft 的微观机制之前,有必要先理解 CP 协议在 Nacos 集群一致性体系中的定位,以及它与 AP 协议(Distro)的职责边界。
1.1 AP 与 CP 的职责边界
Nacos 同时支持 AP 和 CP 两种一致性模式,分别服务于不同类型的数据:
AP 模式(Distro 协议):服务注册发现的临时实例(ephemeral=true),采用最终一致性模型。这类数据的特点是生命周期短、变更频繁、允许短时间不一致,最终通过定时校验收敛。我们在第十二至十四篇文章中已完整分析了 Distro 协议的全部机制。
CP 模式(JRaft 协议):持久化实例(ephemeral=false)和配置数据,采用强一致性模型。这类数据的特点是持久存储、不允许丢数据、读取必须返回最新值。Nacos 在 2.x 版本中引入 SOFAJRaft 作为 Raft 协议的实现,为 CP 模式提供 Leader 选举、日志复制、快照管理能力。
两种协议在启动阶段通过 ProtocolManager 统一管理,互不干扰。同一 Nacos 节点可以同时运行 Distro 和 JRaft 协议——前者负责临时实例的最终一致性,后者负责持久化数据和配置的强一致性。
1.2 一致性协议 SPI 层
CP 协议遵循一套 SPI 接口体系,定义在 consistency 模块中:
ConsistencyProtocol 接口定义了协议的通用契约:init() 初始化协议配置,getData()/aGetData() 提供同步/异步读取,write()/writeAsync() 提供同步/异步写入,memberChange() 处理集群成员变更,protocolMetaData() 暴露协议运行时元数据(如当前 Leader 地址和任期)。
CPProtocol 接口在 ConsistencyProtocol 基础上扩展了一个关键方法 isLeader(String group)。这个加在接口层面的方法体现了 CP 模式的核心需求——业务方必须能够判断当前节点对于某个 Raft Group 是否为 Leader。这个判断在读写请求的转发决策中起到关键作用。
RequestProcessor4CP 是业务方需要实现的抽象类,定义了四个方法:group() 返回该 Processor 所属的 Raft Group 名称;onApply(WriteRequest) 是 Raft 日志提交后的状态机回放入口;onRequest(ReadRequest) 处理线性一致性读取;loadSnapshotOperate() 返回快照操作列表。
二、JRaft 整体架构
JRaft 协议的架构分为三个层次,核心是 JRaftProtocol 和 JRaftServer 的配合。本节的时序图展示了从 Spring 启动到 Raft Group 完全就绪的初始化链路:

2.1 ProtocolManager:协议注册与成员变更通知
ProtocolManager 是 CP 和 AP 协议的统一管理入口,使用 Spring @Component 标注,在容器启动时被自动装配:
com.alibaba.nacos.core.distributed.ProtocolManager#getCpProtocol
public CPProtocol getCpProtocol() {
if (!cpInit){
synchronized (cpLock) {
if (!cpInit) {
initCPProtocol();
cpInit = true;
}
}
}
return cpProtocol;
}
initCPProtocol() 通过 Spring 获取 CPProtocol 类型的 Bean,该 Bean 由 ConsistencyConfiguration 注册:
com.alibaba.nacos.core.distributed.ConsistencyConfiguration#strongAgreementProtocol
@Bean(value = "strongAgreementProtocol")
public CPProtocol strongAgreementProtocol(ServerMemberManager memberManager) throws Exception {
final CPProtocol protocol = getProtocol(CPProtocol.class,
() -> new JRaftProtocol(memberManager));
return protocol;
}
这里优先通过 SPI 加载 CPProtocol 实现(允许用户自定义),若未找到则使用默认的 JRaftProtocol。initCPProtocol() 获取 Bean 后,通过 injectMembers4CP(config) 将当前集群成员注入 Raft 配置,然后调用 protocol.init(config) 启动 JRaft 协议。
ProtocolManager 还实现了 MemberChangeListener,在集群成员变更时同时更新 AP 和 CP 协议:
com.alibaba.nacos.core.distributed.ProtocolManager#onEvent
@Override
public void onEvent(MembersChangeEvent event) {
if (Objects.nonNull(apProtocol)) {
ProtocolExecutor.apMemberChange(() -> apProtocol.memberChange(toAPMembersInfo(event.getMembers())));
}
if (Objects.nonNull(cpProtocol)) {
ProtocolExecutor.cpMemberChange(() -> cpProtocol.memberChange(toCPMembersInfo(event.getMembers())));
}
}
2.2 JRaftProtocol.init():Raft 协议初始化
JRaftProtocol.init() 分三个阶段完成初始化:
com.alibaba.nacos.core.distributed.raft.JRaftProtocol#init
@Override
public void init(RaftConfig config) {
if (initialized.compareAndSet(false, true)) {
this.raftConfig = config;
NotifyCenter.registerToSharePublisher(RaftEvent.class);
this.raftServer.init(this.raftConfig);
this.raftServer.start();
NotifyCenter.registerSubscriber(new Subscriber<RaftEvent>() {
@Override
public void onEvent(RaftEvent event) {
final String groupId = event.getGroupId();
Map<String, Object> properties = new HashMap<>();
MapUtil.putIfValNoEmpty(properties, MetadataKey.LEADER_META_DATA, event.getLeader());
MapUtil.putIfValNoNull(properties, MetadataKey.TERM_META_DATA, event.getTerm());
MapUtil.putIfValNoEmpty(properties, MetadataKey.RAFT_GROUP_MEMBER, event.getRaftClusterInfo());
MapUtil.putIfValNoEmpty(properties, MetadataKey.ERR_MSG, event.getErrMsg());
value.put(groupId, properties);
metaData.load(value);
injectProtocolMetaData(metaData);
}
});
}
}
第一阶段注册 RaftEvent 到共享事件发布器。第二阶段初始化并启动 JRaftServer(解析本机地址、初始化 NodeOptions、创建 RpcServer)。第三阶段注册 RaftEvent 订阅者——每当 Leader 变更时,该订阅者将最新的 leader、term、raftGroupMember 等元数据注入 ProtocolMetaData 中,供业务方通过 protocol.protocolMetaData() 查询。
三、Multi-Raft Group 机制
JRaft 协议的一个核心设计是Multi-Raft Group(多 Raft 组)。与大多数仅维护一个 Raft 组的应用不同,Nacos 在同一进程中运行多个独立的 Raft 组,各自维护自己的状态机、选举和日志复制。
3.1 为什么需要多个 Raft Group
JRaftServer 类注释给出了清晰的解释:
每个 LogProcessor 对应不同的功能模块,如 Nacos 的 naming 模块和 config 模块。这两个模块相互独立、互不影响。如果只有一个状态机,所有功能模块的日志处理就会混杂在一起,任何一个模块在日志处理过程中出现异常或长时间阻塞,都会影响其他功能模块的正常运行。
Nacos 2.4.3 中定义了如下 Raft Group:
| Raft Group 名称 | Processor | 所属模块 | 用途 |
|---|---|---|---|
naming_persistent_service_v2 |
PersistentClientOperationServiceImpl |
Naming | 持久化实例注册/注销 |
naming_service_metadata |
ServiceMetadataProcessor |
Naming | Service 元数据 |
naming_instance_metadata |
InstanceMetadataProcessor |
Naming | Instance 元数据 |
nacos_config |
DistributedDatabaseOperateImpl |
Config | Derby 数据库 SQL 日志 |
每个 Raft Group 独立运行,拥有自己的 Leader、日志序列和快照周期,互不干扰。
下面的对比图直观展示了单 Raft Group 与 Multi-Raft Group 在模块隔离、日志独立性和故障隔离上的本质差异:

3.2 createMultiRaftGroup():为每个 Processor 创建 Raft Group
当 RequestProcessor4CP 的实现类通过 protocol.addRequestProcessors(Collections.singletonList(this)) 向 JRaft 注册时,JRaftProtocol 将注册委托给 JRaftServer.createMultiRaftGroup():
com.alibaba.nacos.core.distributed.raft.JRaftServer#createMultiRaftGroup
synchronized void createMultiRaftGroup(Collection<RequestProcessor4CP> processors) {
if (!this.isStarted) {
this.processors.addAll(processors);
return;
}
final String parentPath = Paths.get(EnvUtil.getNacosHome(), "data/protocol/raft").toString();
for (RequestProcessor4CP processor : processors) {
final String groupName = processor.group();
if (multiRaftGroup.containsKey(groupName)) {
throw new DuplicateRaftGroupException(groupName);
}
Configuration configuration = conf.copy();
NodeOptions copy = nodeOptions.copy();
JRaftUtils.initDirectory(parentPath, groupName, copy);
NacosStateMachine machine = new NacosStateMachine(this, processor);
copy.setFsm(machine);
copy.setInitialConf(configuration);
int doSnapshotInterval = ConvertUtils.toInt(
raftConfig.getVal(RaftSysConstants.RAFT_SNAPSHOT_INTERVAL_SECS),
RaftSysConstants.DEFAULT_RAFT_SNAPSHOT_INTERVAL_SECS);
doSnapshotInterval = CollectionUtils.isEmpty(processor.loadSnapshotOperate()) ? 0 : doSnapshotInterval;
copy.setSnapshotIntervalSecs(doSnapshotInterval);
RaftGroupService raftGroupService = new RaftGroupService(groupName, localPeerId, copy, rpcServer, true);
Node node = raftGroupService.start(false);
machine.setNode(node);
RouteTable.getInstance().updateConfiguration(groupName, configuration);
RaftExecutor.executeByCommon(() -> registerSelfToCluster(groupName, localPeerId, configuration));
Random random = new Random();
long period = nodeOptions.getElectionTimeoutMs() + random.nextInt(5 * 1000);
RaftExecutor.scheduleRaftMemberRefreshJob(() -> refreshRouteTable(groupName),
nodeOptions.getElectionTimeoutMs(), period, TimeUnit.MILLISECONDS);
multiRaftGroup.put(groupName, new RaftGroupTuple(node, processor, raftGroupService, machine));
}
}
创建过程的几个关键步骤:
配置隔离。每个 Raft Group 拥有独立的 Configuration 和 NodeOptions 拷贝。conf.copy() 和 nodeOptions.copy() 保证了基础配置共享,修改不互相污染。快照目录通过 JRaftUtils.initDirectory(parentPath, groupName, copy) 生成独立的日志和快照存储路径(如 data/protocol/raft/naming_persistent_service_v2/)。
状态机关联。每个 Raft Group 绑定一个 NacosStateMachine 实例,该实例在 onApply() 中委托给对应的 RequestProcessor4CP.onApply()。这意味着 naming 持久化实例的日志提交只影响 naming_persistent_service_v2 组的状态机,不会干扰 nacos_config 组的状态机。
Leader 刷新任务。创建完成后为每个组启动一个周期性任务 refreshRouteTable(groupName),以 electionTimeoutMs + random(5s) 为周期刷新该组的 Leader 和路由节点配置。刷新周期的随机偏移量避免了所有组同时刷新导致的计算尖峰。
3.3 NacosStateMachine:业务状态机
NacosStateMachine 继承自 SOFAJRaft 的 StateMachineAdapter,是连接 JRaft 底层 Raft 算法与 Nacos 业务逻辑的桥梁:
com.alibaba.nacos.core.distributed.raft.NacosStateMachine
class NacosStateMachine extends StateMachineAdapter {
protected final JRaftServer server;
protected final RequestProcessor processor;
private final AtomicBoolean isLeader = new AtomicBoolean(false);
private final String groupId;
private volatile long term = -1;
private volatile String leaderIp = "unknown";
private Node node;
@Override
public void onApply(Iterator iter) {
while (iter.hasNext()) {
Message message;
if (iter.done() != null) {
closure = (NacosClosure) iter.done();
message = closure.getMessage();
} else {
final ByteBuffer data = iter.getData();
message = ProtoMessageUtil.parse(data.array());
if (message instanceof ReadRequest) {
applied++; iter.next(); continue;
}
}
if (message instanceof WriteRequest) {
Response response = processor.onApply((WriteRequest) message);
postProcessor(response, closure);
}
if (message instanceof ReadRequest) {
Response response = processor.onRequest((ReadRequest) message);
postProcessor(response, closure);
}
closure.run(status);
applied++; iter.next();
}
}
@Override
public void onLeaderStart(final long term) {
super.onLeaderStart(term);
this.term = term;
this.isLeader.set(true);
this.leaderIp = node.getNodeId().getPeerId().getEndpoint().toString();
NotifyCenter.publishEvent(
RaftEvent.builder().groupId(groupId).leader(leaderIp).term(term)
.raftClusterInfo(allPeers()).build());
}
@Override
public void onLeaderStop(final Status status) {
super.onLeaderStop(status);
this.isLeader.set(false);
}
}
onApply() 在 Raft 日志提交后被调用。当 iter.done() != null 时,表示是 Leader 节点提交的本地任务(NacosClosure 携带了原始消息和回调),直接从 closure 中获取消息;iter.done() == null 时,表示是 Follower 节点从 Leader 复制过来的日志,需要解析 ByteBuffer 获取消息。如果是 WriteRequest,委托给 processor.onApply() 执行业务逻辑;如果是 ReadRequest,Follower 节点会跳过读取请求(continue),因为只有 Leader 能在日志应用后执行线性一致性读取。
四、Leader 选举机制
Leader 选举是 Raft 协议最核心的机制。在 Nacos 的实现中,JRaft 底层选举进程由 SOFAJRaft 框架处理,但 Nacos 通过配置参数、状态机回调和事件发布参与了选举的全生命周期。
4.1 选举配置参数
Leader 选举相关的配置集中在 JRaftServer.init() 中初始化:
com.alibaba.nacos.core.distributed.raft.JRaftServer#init
int electionTimeout = Math.max(
ConvertUtils.toInt(config.getVal(RaftSysConstants.RAFT_ELECTION_TIMEOUT_MS),
RaftSysConstants.DEFAULT_ELECTION_TIMEOUT),
RaftSysConstants.DEFAULT_ELECTION_TIMEOUT);
nodeOptions.setSharedElectionTimer(true);
nodeOptions.setSharedVoteTimer(true);
nodeOptions.setSharedStepDownTimer(true);
nodeOptions.setSharedSnapshotTimer(true);
nodeOptions.setElectionTimeoutMs(electionTimeout);
选举超时(electionTimeout):默认 5000 毫秒,通过 nacos.core.protocol.raft.election_timeout_ms 配置。该值决定了 Follower 在未收到 Leader 心跳后等待多久才转为 Candidate 发起选举。Math.max(..., DEFAULT_ELECTION_TIMEOUT) 保证了即使配置了更小的值也会强制提升到 5000ms,防止因超时过短导致频繁选举。
共享定时器:setSharedElectionTimer(true) 和 setSharedVoteTimer(true) 让所有 Raft Group 的选举定时器和投票定时器共享同一个线程池。这是因为在 Multi-Raft Group 模式下,每个 Group 都需要自己的定时器来触发选举超时,如果每个 Group 使用独立的线程池,Group 数量较多时会占用大量线程资源。共享定时器通过一个 SharedTimer 管理所有 Group 的时间事件,大幅降低线程数。
心跳与超时关系:DEFAULT_ELECTION_HEARTBEAT_FACTOR = 10 定义了选举超时与心跳间隔的比值。心跳间隔 = electionTimeoutMs / electionHeartbeatFactor = 5000 / 10 = 500ms。即 Leader 每 500ms 向 Follower 发送一次心跳续约,Follower 在 5000ms 内未收到心跳则触发选举。
随机化防冲突:DEFAULT_MAX_ELECTION_DELAY_MS = 1000 定义了选举超时的随机偏移量。Follower 的实际选举超时时间在 [electionTimeout, electionTimeout + maxElectionDelay] 之间随机选择。这意味着 3 节点集群中,三个 Follower 各自的超时时间可能在 [5000ms, 6000ms] 之间随机分布,避免了所有 Follower 同时超时引发冲突选举。
4.2 onLeaderStart/onLeaderStop 回调
当 SOFAJRaft 底层的 Raft 算法完成投票并选出 Leader 后,会调用 NacosStateMachine.onLeaderStart(term):
com.alibaba.nacos.core.distributed.raft.NacosStateMachine#onLeaderStart
@Override
public void onLeaderStart(final long term) {
super.onLeaderStart(term);
this.term = term;
this.isLeader.set(true);
this.leaderIp = node.getNodeId().getPeerId().getEndpoint().toString();
NotifyCenter.publishEvent(
RaftEvent.builder()
.groupId(groupId)
.leader(leaderIp)
.term(term)
.raftClusterInfo(allPeers())
.build());
}
onLeaderStart 执行三个动作:更新内部状态(isLeader=true、记录当前 term、记录 leaderIp 为本机地址),发布 RaftEvent。事件包含 groupId(哪个 Raft Group 产生了新 Leader)、leader(Leader 地址)、term(当前任期号)和 raftClusterInfo(当前 Group 的所有节点列表)。
onLeaderStop 的触发场景包括 Leader 节点崩溃、网络分区导致新 Leader 被选出、或者 Leader 主动退位:
com.alibaba.nacos.core.distributed.raft.NacosStateMachine#onLeaderStop
@Override
public void onLeaderStop(final Status status) {
super.onLeaderStop(status);
this.isLeader.set(false);
}
onLeaderStop 仅做一件事——将 isLeader 置为 false。它不发布事件,也不清理 leaderIp 和 term。这是因为 onLeaderStop 可能伴随节点崩溃,发布事件本身可能失败。而 JRaftProtocol 注册的 RaftEvent 订阅者会在下一次 onLeaderStart 时用新值覆盖元数据。
4.3 RaftEvent 的元数据更新闭环
JRaftProtocol.init() 中注册的 RaftEvent 订阅者,将事件信息注入 ProtocolMetaData:
com.alibaba.nacos.core.distributed.raft.JRaftProtocol
NotifyCenter.registerSubscriber(new Subscriber<RaftEvent>() {
@Override
public void onEvent(RaftEvent event) {
final String groupId = event.getGroupId();
Map<String, Object> properties = new HashMap<>();
MapUtil.putIfValNoEmpty(properties, MetadataKey.LEADER_META_DATA, event.getLeader());
MapUtil.putIfValNoNull(properties, MetadataKey.TERM_META_DATA, event.getTerm());
MapUtil.putIfValNoEmpty(properties, MetadataKey.RAFT_GROUP_MEMBER, event.getRaftClusterInfo());
value.put(groupId, properties);
metaData.load(value);
injectProtocolMetaData(metaData);
}
});
ProtocolMetaData 是 CP 协议暴露给外部使用的运行时元数据集,包含每个 Raft Group 的 Leader 地址、任期和集群成员信息。这些元数据通过 CPProtocol.protocolMetaData() 暴露给业务方(如 Nacos 控制台),用于判断当前集群中每个 Raft Group 的 Leader 节点,以及观测是否存在 Leader 缺失的情况。
订阅者的 MapUtil.putIfValNoEmpty / putIfValNoNull 保证了空值不会污染元数据。例如当 Raft Group 尚未完成选举(无 Leader)时,event.getLeader() 为空字符串,不会被写入元数据。
下面的闭环示意图汇总了从心跳超时触发选举到路由就绪的完整生命周期:

4.4 isReady() 严格模式
JRaftServer.isReady() 用于判断当前节点是否已完全就绪,分为两种模式:
com.alibaba.nacos.core.distributed.raft.JRaftServer#isReady
public boolean isReady() {
if (raftConfig.isStrictMode()) {
for (RequestProcessor4CP each : processors) {
if (null == getLeader(each.group())) {
return false;
}
}
}
return isStarted;
}
严格模式下,只有当所有 Raft Group 都有 Leader 时 isReady() 才返回 true。这避免了"节点启动但某些 Group 的 Leader 尚未选举完成就对外提供服务"的问题。非严格模式下,仅检查 isStarted(Raft 协议是否已启动)。
五、读写请求的路由策略
选举出 Leader 后,业务方通过 JRaftProtocol 发起读写请求。请求的路由策略取决于当前节点在目标 Raft Group 中的角色。
下面的决策树展示了写入请求根据当前节点角色(Leader / Follower / Candidate)分别走到的三条处理路径:

5.1 写入请求的路由
写入请求的入口是 JRaftProtocol.write(),它调用 writeAsync() 将请求委托给 JRaftServer.commit():
com.alibaba.nacos.core.distributed.raft.JRaftProtocol#writeAsync
@Override
public CompletableFuture<Response> writeAsync(WriteRequest request) {
return raftServer.commit(request.getGroup(), request, new CompletableFuture<>());
}
commit() 方法根据当前节点的角色分流:
com.alibaba.nacos.core.distributed.raft.JRaftServer#commit
public CompletableFuture<Response> commit(final String group, final Message data,
final CompletableFuture<Response> future) {
final RaftGroupTuple tuple = findTupleByGroup(group);
if (tuple == null) {
future.completeExceptionally(
new IllegalArgumentException("No corresponding Raft Group found : " + group));
return future;
}
FailoverClosureImpl closure = new FailoverClosureImpl(future);
final Node node = tuple.node;
if (node.isLeader()) {
applyOperation(node, data, closure);
} else {
invokeToLeader(group, data, rpcRequestTimeoutMs, closure);
}
return future;
}
Leader 路径:当前节点是该 Raft Group 的 Leader 时,直接调用 applyOperation() 将请求作为 Task 提交到 Raft 日志:
com.alibaba.nacos.core.distributed.raft.JRaftServer#applyOperation
public void applyOperation(Node node, Message data, FailoverClosure closure) {
final Task task = new Task();
task.setDone(new NacosClosure(data, status -> { ... }));
// 在头部添加请求类型字段,onApply 时根据该字段区分 Read/Write
byte[] requestTypeFieldBytes = new byte[2];
requestTypeFieldBytes[0] = ProtoMessageUtil.REQUEST_TYPE_FIELD_TAG;
if (data instanceof ReadRequest) {
requestTypeFieldBytes[1] = ProtoMessageUtil.REQUEST_TYPE_READ;
} else {
requestTypeFieldBytes[1] = ProtoMessageUtil.REQUEST_TYPE_WRITE;
}
byte[] dataBytes = data.toByteArray();
task.setData((ByteBuffer) ByteBuffer.allocate(requestTypeFieldBytes.length + dataBytes.length)
.put(requestTypeFieldBytes).put(dataBytes).position(0));
node.apply(task);
}
node.apply(task) 将请求提交到 JRaft 的日志写入队列,写入本地日志文件后,通过 Raft 复制到 Follower 节点。当日志被大多数节点(Quorum)确认后,NacosStateMachine.onApply() 被调用执行状态机更新,并通过 NacosClosure 回调返回结果给调用方。
Follower 路径:当前节点是 Follower 时,调用 invokeToLeader() 将请求转发给 Leader:
com.alibaba.nacos.core.distributed.raft.JRaftServer#invokeToLeader
private void invokeToLeader(final String group, final Message request, final int timeoutMillis,
FailoverClosure closure) {
try {
final Endpoint leaderIp = Optional.ofNullable(getLeader(group))
.orElseThrow(() -> new NoLeaderException(group)).getEndpoint();
cliClientService.getRpcClient().invokeAsync(leaderIp, request, new InvokeCallback() {
@Override
public void complete(Object o, Throwable ex) {
if (Objects.nonNull(ex)) {
closure.setThrowable(ex);
closure.run(new Status(RaftError.UNKNOWN, ex.getMessage()));
return;
}
if (!((Response) o).getSuccess()) {
closure.setThrowable(new IllegalStateException(((Response) o).getErrMsg()));
closure.run(new Status(RaftError.UNKNOWN, ((Response) o).getErrMsg()));
return;
}
closure.setResponse((Response) o);
closure.run(Status.OK());
}
@Override
public Executor executor() {
return RaftExecutor.getRaftCliServiceExecutor();
}
}, timeoutMillis);
} catch (Exception e) {
closure.setThrowable(e);
closure.run(new Status(RaftError.UNKNOWN, e.toString()));
}
}
Follower 首先通过 getLeader(group) 查询 RouteTable 获取该 Group 当前 Leader 的 Endpoint。如果找不到 Leader(如选举尚未完成),抛出 NoLeaderException。找到 Leader 后,通过 cliClientService.getRpcClient().invokeAsync() 发送异步 RPC 请求到 Leader。超时时间由 rpcRequestTimeoutMs 控制(默认 5000ms)。
NoLeaderException 是 JRaft 模块自定义的异常,调用方需要捕获并处理(如重试或返回错误给客户端):
com.alibaba.nacos.core.distributed.raft.exception.NoLeaderException
public class NoLeaderException extends Exception {
public NoLeaderException(String group) {
super("The Raft Group [" + group + "] did not find the Leader node");
}
}
5.2 读取请求与线性一致性
与写入请求不同,读取请求需要通过 JRaftServer.get() 遵循 Raft 的线性一致性读取协议:
com.alibaba.nacos.core.distributed.raft.JRaftServer#get
CompletableFuture<Response> get(final ReadRequest request) {
final String group = request.getGroup();
CompletableFuture<Response> future = new CompletableFuture<>();
final RaftGroupTuple tuple = findTupleByGroup(group);
if (Objects.isNull(tuple)) {
future.completeExceptionally(new NoSuchRaftGroupException(group));
return future;
}
final Node node = tuple.node;
final RequestProcessor processor = tuple.processor;
try {
node.readIndex(BytesUtil.EMPTY_BYTES, new ReadIndexClosure() {
@Override
public void run(Status status, long index, byte[] reqCtx) {
if (status.isOk()) {
try {
Response response = processor.onRequest(request);
future.complete(response);
} catch (Throwable t) {
future.completeExceptionally(new ConsistencyException(...));
}
return;
}
readFromLeader(request, future);
}
});
return future;
} catch (Throwable e) {
readFromLeader(request, future);
return future;
}
}
线性一致性读取分两个阶段:第一阶段通过 node.readIndex() 向 Leader 确认当前 Raft Group 的"已提交日志索引"。Leader 返回最新已提交的日志索引后,第二阶段直接在本地状态机执行 processor.onRequest(request) 读取数据。这种"先确认已提交索引再读取"的机制保证了读取操作返回的是至少包含所有已提交写入的最新数据视图,实现了强一致性。
如果 readIndex() 返回的 status 不是 isOk()(如当前节点尚未与 Leader 完成同步),或 readIndex() 调用本身抛出异常,则回退到 readFromLeader()——通过写入路径将读取请求转发到 Leader 执行。
在 NacosStateMachine.onApply() 中,Follower 节点会跳过 ReadRequest:
if (iter.done() == null) {
final ByteBuffer data = iter.getData();
message = ProtoMessageUtil.parse(data.array());
if (message instanceof ReadRequest) {
applied++; index++; iter.next();
continue;
}
}
这是因为 ReadIndex 协议下,读取操作不需要通过日志复制到 Follower。Follower 节点虽然也会收到 Leader 的读取日志复制,但 iter.done() == null 的判断使其跳过处理,避免在 Follower 上重复执行读取操作。
六、成员变更与节点自注册
集群成员变更是 Raft 协议中另一个重要环节,涉及到节点加入、退出时 Raft Group 配置的更新。
6.1 集群成员变更
当 ProtocolManager 通过 MemberChangeListener 收到集群成员变更事件时,会调用 JRaftProtocol.memberChange():
com.alibaba.nacos.core.distributed.raft.JRaftProtocol#memberChange
@Override
public void memberChange(Set<String> addresses) {
for (int i = 0; i < 5; i++) {
if (this.raftServer.peerChange(jRaftMaintainService, addresses)) {
return;
}
ThreadUtils.sleep(100L);
}
Loggers.RAFT.warn("peer removal failed");
}
最多重试 5 次,每次间隔 100ms。peerChange() 的实现专注于处理节点移除——Raft 协议中节点加入是自发的(registerSelfToCluster),但节点退出需要 Leader 主动将退出的节点从配置中移除:
com.alibaba.nacos.core.distributed.raft.JRaftServer#peerChange
boolean peerChange(JRaftMaintainService maintainService, Set<String> newPeers) {
Set<String> oldPeers = new HashSet<>(this.raftConfig.getMembers());
this.raftConfig.setMembers(localPeerId.toString(), newPeers);
oldPeers.removeAll(newPeers);
if (oldPeers.isEmpty()) { return true; }
AtomicInteger successCnt = new AtomicInteger(0);
multiRaftGroup.forEach((group, tuple) -> {
Map<String, String> params = new HashMap<>();
params.put(JRaftConstants.GROUP_ID, group);
params.put(JRaftConstants.COMMAND_NAME, JRaftConstants.REMOVE_PEERS);
params.put(JRaftConstants.COMMAND_VALUE, StringUtils.join(oldPeers, StringUtils.COMMA));
RestResult<String> result = maintainService.execute(params);
if (result.ok()) { successCnt.incrementAndGet(); }
});
return successCnt.get() == multiRaftGroup.size();
}
遍历所有 Raft Group,对每个组执行 REMOVE_PEERS 操作。只有当所有组的移除操作都成功时,peerChange() 才返回 true。oldPeers 计算为旧成员减去新成员的差值,即需要移除的节点。
6.2 节点自注册
新节点启动时,createMultiRaftGroup() 异步提交 registerSelfToCluster() 任务:
com.alibaba.nacos.core.distributed.raft.JRaftServer#registerSelfToCluster
void registerSelfToCluster(String groupId, PeerId selfIp, Configuration conf) {
while (!isShutdown) {
try {
List<PeerId> peerIds = cliService.getPeers(groupId, conf);
if (peerIds.contains(selfIp)) {
return;
}
Status status = cliService.addPeer(groupId, conf, selfIp);
if (status.isOk()) {
return;
}
Loggers.RAFT.warn("Failed to join the cluster, retry...");
} catch (Exception e) {
Loggers.RAFT.error("Failed to join the cluster, retry...", e);
}
ThreadUtils.sleep(1_000L);
}
}
首先通过 cliService.getPeers() 检查当前节点是否已在 Raft 集群中(用于节点重启场景)。如果已在,直接返回。如果不在,调用 cliService.addPeer() 将自身加入集群。失败时以 1 秒为间隔持续重试,直到成功或 isShutdown。
这个"检视已存在"的设计支持了滚动重启场景——节点重启后检测到自身已在集群配置中,直接跳过注册步骤,Raft Group 可以继续正常运作。
全文小结
本文聚焦 JRaft 架构的 CP 一致性实现,从整体架构、Multi-Raft Group 机制、Leader 选举到读写请求路由,完整分析了 Nacos 中 CP 模式的核心运行机制。
在整体架构方面,ProtocolManager 负责 AP 和 CP 协议的统一管理,JRaftProtocol 通过 SPI 机制加载并初始化为 Raft 协议的实现入口,JRaftServer 是底层 Raft 引擎的载体。在 Multi-Raft Group 方面,每个 RequestProcessor4CP 对应一个独立的 Raft Group,拥有自己的 NacosStateMachine、日志序列和快照周期,保证了不同功能模块(命名服务、配置服务、元数据服务)的 Raft 运行互不干扰。在 Leader 选举方面,JRaftServer.init() 通过 electionTimeout(5000ms)、共享定时器、随机化防冲突等配置参数为选举提供基础条件,NacosStateMachine 的 onLeaderStart/onLeaderStop 回调发布 RaftEvent,由 JRaftProtocol 的订阅者更新 ProtocolMetaData,形成选举到元数据感知的完整闭环。在读写路由方面,commit() 根据 node.isLeader() 分流——Leader 直接提交 Raft 日志,Follower 通过 invokeToLeader 转发给 Leader;get() 通过 readIndex 实现线性一致性读取,确保读操作不落后于已提交的写入。
原创不易,如果本文对您有帮助,带来了些许灵感或启发,烦请动动小手点赞、关注、转发、收藏。这是作者持续更新的动力源泉,衷心感谢您的支持。我会尽量在工作之余,为大家带来更高品质的内容,努力保持周更。
更多推荐




所有评论(0)