在这里插入图片描述

《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 协议异步数据同步的完整内部机制,理解了延迟任务合并、双模式调度、传输层前置检查、增量同步算法和启动加载五大核心模块。

但 Distro 协议如何发现异步同步可能遗漏的数据差异?DistroVerifyTimedTask 的周期性校验如何在不打断正常业务的情况下,逐步收敛集群数据到一致状态?接收端 processVerifyData 的 revision 比较算法如何精确识别数据版本差异,而 ClientVerifyFailedEvent 又是如何触发冲突修复,将正确的客户端数据重新同步到校验失败的节点?本文将聚焦 Distro 协议的容错处理机制,从校验任务的生命周期、执行逻辑、验证算法到冲突修复,完整剖析定时校验如何兜底保证 AP 模式的最终一致性。

一、数据校验机制全景回顾

1.1 校验机制在 Distro 协议中的定位

在上一篇文章中,我们提到 Distro 协议的第三个设计思想是"增量同步 + 校验兜底"。增量同步负责在事件驱动的异步链路中高效地传输变更数据,但单次同步可能因网络故障、节点宕机或消息丢失而遗漏数据。校验兜底作为最后的防线,通过周期性轮询而非事件驱动的方式,主动发现集群节点间的数据差异并触发修复,保证数据在长周期下最终收敛到一致状态。

校验机制与主动同步有三个本质区别:

一是驱动方式不同。主动同步是事件驱动的(客户端变更触发 ClientChangedEvent),只要有变更就会立即发起同步;校验是定时驱动的(DistroVerifyTimedTask 每 5 秒触发一次),无论有无变更都会周期性执行业务。

二是数据负载不同。主动同步传输的是全量 ClientSyncData(包含客户端的所有服务实例发布信息),数据量大;校验传输的是轻量级的 DistroClientVerifyInfo(仅含 clientId + revision 两个字段),数据量极小。

三是修复策略不同。主动同步失败后会通过 DistroClientTaskFailedHandler 无限重试;校验发现不一致后,不重试验证而是直接触发补偿同步——一步到位发送全量数据到目标节点。这种"发现问题就一次性修复"的策略,避免了校验与重试的循环叠加。

下面的对比图直观展示了主动同步与校验同步在驱动方式、数据负载和修复策略上的三个本质区别:

在这里插入图片描述

1.2 校验生命周期中的关键类和实体

在深入分析之前,先快速了解本文涉及的五个核心类:

  • DistroVerifyTimedTask:周期性校验任务的调度器,实现 Runnable 接口,由 GlobalExecutor.schedulePartitionDataTimedSync() 以 5 秒为周期执行。它的职责是收集各数据存储类型的校验数据,分发给所有非自身成员节点。
  • DistroVerifyExecuteTask:单次校验任务的执行单元,继承 AbstractExecuteTask,负责将一组 DistroData 校验数据逐个发送到目标节点,根据传输代理是否支持回调分流。
  • DistroClientVerifyInfo:校验数据的最小载体,仅包含 clientId(客户端标识)和 revision(数据版本号)两个字段。这个极简的数据结构是整个校验协议的基础——只需两个字段就能完成数据一致性的判定。
  • DistroVerifyCallbackWrapper:异步回调的传输层实现,在校验失败时负责发布 ClientVerifyFailedEvent 和更新 TPS 监控指标。
  • ClientVerifyFailedEvent:校验失败事件,包含 clientIdtargetServer 两个字段,监听者是 DistroClientDataProcessor(它同时是 SmartSubscriber),触发补偿同步。

二、DistroVerifyTimedTask:校验任务的生命周期

校验任务的生命周期从节点启动开始,贯穿集群运行的整个周期。本节的时序图展示了从 startDistroTask()run() 执行、再到与数据存储层交互的完整调度链路:

在这里插入图片描述

2.1 启动时机与调度策略

校验任务的启动与加载任务(DistroLoadDataTask)同时创建于 DistroProtocol.startDistroTask() 方法中:

com.alibaba.nacos.core.distributed.distro.DistroProtocol#startDistroTask

private void startDistroTask() {
    if (EnvUtil.getStandaloneMode()) {
        isInitialized = true;
        return;
    }
    startVerifyTask();
    startLoadTask();
}

单机模式下,isInitialized 直接置为 true,跳过所有校验和加载任务——因为单机无需数据同步,自然也不需校验。集群模式下,校验任务和加载任务先后启动。

com.alibaba.nacos.core.distributed.distro.DistroProtocol#startVerifyTask

private void startVerifyTask() {
    GlobalExecutor.schedulePartitionDataTimedSync(
            new DistroVerifyTimedTask(memberManager, distroComponentHolder,
                    distroTaskEngineHolder.getExecuteWorkersManager()),
            DistroConfig.getInstance().getVerifyIntervalMillis());
}

校验任务通过 GlobalExecutor.schedulePartitionDataTimedSync() 进行周期性调度,默认周期为 5000 毫秒(5 秒),可通过配置项 nacos.core.protocol.distro.data.verify.intervalMs 调整。5 秒的间隔考虑了校验的及时性与网络开销的平衡——太频繁(如 1 秒)会增加网络负载,太低频(如 30 秒)则数据不一致的时间窗口过长。


2.2 run() 方法的三级执行流程

DistroVerifyTimedTask.run() 的执行逻辑可以概括为"获取节点列表 → 遍历存储类型 → 逐类型分发"的三级流程:

com.alibaba.nacos.core.distributed.distro.task.verify.DistroVerifyTimedTask#run

@Override
public void run() {
    try {
        List<Member> targetServer = serverMemberManager.allMembersWithoutSelf();
        if (Loggers.DISTRO.isDebugEnabled()) {
            Loggers.DISTRO.debug("server list is: {}", targetServer);
        }
        for (String each : distroComponentHolder.getDataStorageTypes()) {
            verifyForDataStorage(each, targetServer);
        }
    } catch (Exception e) {
        Loggers.DISTRO.error("[DISTRO-FAILED] verify task failed.", e);
    }
}

第一级获取非自身成员的目标节点列表。第二级从 DistroComponentHolder 中获取所有已注册的数据存储类型(如 Nacos:Naming:v2:ClientData)。第三级为每个类型执行 verifyForDataStorage()run() 的最外层用 try-catch 包裹,防止某个类型的校验异常阻断后续类型的校验——即使 ClientData 的校验抛出异常,DistroRecords 或其他数据类型的校验也能正常进行。


2.3 verifyForDataStorage() 的三项检查

verifyForDataStorage() 在执行校验前会经过三项检查:

com.alibaba.nacos.core.distributed.distro.task.verify.DistroVerifyTimedTask#verifyForDataStorage

private void verifyForDataStorage(String type, List<Member> targetServer) {
    DistroDataStorage dataStorage = distroComponentHolder.findDataStorage(type);
    if (!dataStorage.isFinishInitial()) {
        Loggers.DISTRO.warn("data storage {} has not finished initial step, do not send verify data",
                dataStorage.getClass().getSimpleName());
        return;
    }
    List<DistroData> verifyData = dataStorage.getVerifyData();
    if (null == verifyData || verifyData.isEmpty()) {
        return;
    }
    for (Member member : targetServer) {
        DistroTransportAgent agent = distroComponentHolder.findTransportAgent(type);
        if (null == agent) {
            continue;
        }
        executeTaskExecuteEngine.addTask(member.getAddress() + type,
                new DistroVerifyExecuteTask(agent, verifyData, member.getAddress(), type));
    }
}

检查一(初始化状态):dataStorage.isFinishInitial() 检查数据存储是否已完成初始化。在上一篇文章中,我们了解到 DistroClientDataProcessorisFinishInitial 标志在 DistroLoadDataTask 成功加载快照后才被置为 true。未初始化的节点根本不应该发送校验数据——因为它的数据可能本身就少(还未从远端加载完),用它来校验只会产生误报。

检查二(校验数据非空):getVerifyData() 返回当前节点该类型下所有"本节点有责任管理"的客户端的校验数据。如果为空(如节点尚未有客户端连接),直接跳过该校验周期。

检查三(传输代理存在性):findTransportAgent(type) 查找传输代理。如果该类型未注册传输代理(如旧版本未初始化完毕),跳过该节点的校验。

三项检查全部通过后,为目标节点的每类数据类型创建一个 DistroVerifyExecuteTask,以 member.getAddress() + type 作为 Key 提交到执行引擎。Key 的设计保证了同一目标节点同一数据类型的校验任务不会被重复提交(执行引擎的 addTask() 方法是幂等的)。


三、DistroVerifyExecuteTask:校验任务的执行逻辑

执行任务被提交到引擎后,由 DistroVerifyExecuteTask.run() 执行实际的校验数据发送。本节的时序图展示了从执行引擎调度到发送链路、再到收到响应的完整流程:

在这里插入图片描述

3.1 run() 方法的逐条发送策略

com.alibaba.nacos.core.distributed.distro.task.verify.DistroVerifyExecuteTask#run

@Override
public void run() {
    for (DistroData each : verifyData) {
        try {
            if (transportAgent.supportCallbackTransport()) {
                doSyncVerifyDataWithCallback(each);
            } else {
                doSyncVerifyData(each);
            }
        } catch (Exception e) {
            Loggers.DISTRO
                    .error("[DISTRO-FAILED] verify data for type {} to {} failed.", resourceType, targetServer, e);
        }
    }
}

run() 方法逐条遍历 verifyData 列表,对每个校验数据独立发送。这种逐条策略与主动同步的批量合并形成鲜明对比——主动同步需要将同 Key 任务聚合成一次(merge() 算法),而校验任务的每条数据都是独立且互不相干的(每个 clientId 的 revision 独立变更),合并反而会引入错误的语义。

逐条遍历的另一个好处是可靠性——如果某条校验数据发送失败(抛出异常),catch 块记录错误日志后继续处理下一条,不会因为单条失败而阻塞整个校验周期。


3.2 同步与异步的 sendVerifyData 方法

syncVerifyData() 的同步版本包含三道关键逻辑:前置检查、地址替换、发送验证:

com.alibaba.nacos.naming.consistency.ephemeral.distro.v2.DistroClientTransportAgent#syncVerifyData

@Override
public boolean syncVerifyData(DistroData verifyData, String targetServer) {
    if (isNoExistTarget(targetServer)) {
        return true;
    }
    verifyData.getDistroKey().setTargetServer(memberManager.getSelf().getAddress());
    DistroDataRequest request = new DistroDataRequest(verifyData, DataOperation.VERIFY);
    Member member = memberManager.find(targetServer);
    if (checkTargetServerStatusUnhealthy(member)) {
        Loggers.DISTRO
                .warn("[DISTRO] Cancel distro verify caused by target server {} unhealthy, key: {}", targetServer,
                        verifyData.getDistroKey());
        return false;
    }
    try {
        Response response = clusterRpcClientProxy.sendRequest(member, request);
        return checkResponse(response);
    } catch (NacosException e) {
        Loggers.DISTRO.error("[DISTRO-FAILED] Verify distro data failed! key: {} ", verifyData.getDistroKey(), e);
    }
    return false;
}

前置检查的逻辑与主动同步的 syncData() 一致:isNoExistTarget(目标不存在则跳过)和 checkTargetServerStatusUnhealthy(节点不健康则跳过)。与主动同步的一个关键不同是 verifyData.getDistroKey().setTargetServer(memberManager.getSelf().getAddress()) 地址替换——将 DistroKey 中的目标服务器替换为本机地址。这个替换的目的是让接收端在后续的冲突修复(ClientVerifyFailedEvent)中能够正确溯源到"是哪个节点发现数据不一致的",从而触发从源节点到该目标节点的补偿同步。

异步版本的 syncVerifyData() 使用 DistroVerifyCallbackWrapper 作为回调,在校验失败时发布 ClientVerifyFailedEvent

com.alibaba.nacos.naming.consistency.ephemeral.distro.v2.DistroClientTransportAgent#syncVerifyData(3-arg)

@Override
public void syncVerifyData(DistroData verifyData, String targetServer, DistroCallback callback) {
    if (isNoExistTarget(targetServer)) {
        callback.onSuccess();
        return;
    }
    DistroDataRequest request = new DistroDataRequest(verifyData, DataOperation.VERIFY);
    Member member = memberManager.find(targetServer);
    if (checkTargetServerStatusUnhealthy(member)) {
        Loggers.DISTRO.warn("[DISTRO] Cancel distro verify caused by target server {} unhealthy, key: {}",
                targetServer, verifyData.getDistroKey());
        callback.onFailed(null);
        return;
    }
    try {
        DistroVerifyCallbackWrapper wrapper = new DistroVerifyCallbackWrapper(targetServer,
                verifyData.getDistroKey().getResourceKey(), callback, member);
        clusterRpcClientProxy.asyncRequest(member, request, wrapper);
    } catch (NacosException nacosException) {
        callback.onFailed(nacosException);
    }
}

3.3 校验数据的组装

校验数据由 DistroClientDataProcessor.getVerifyData() 方法组装。这个方法遍历所有客户端,收集本节点有责任管理的临时客户端的校验信息:

com.alibaba.nacos.naming.consistency.ephemeral.distro.v2.DistroClientDataProcessor#getVerifyData

@Override
public List<DistroData> getVerifyData() {
    List<DistroData> result = null;
    for (String each : clientManager.allClientId()) {
        Client client = clientManager.getClient(each);
        if (null == client || !client.isEphemeral()) {
            continue;
        }
        if (clientManager.isResponsibleClient(client)) {
            DistroClientVerifyInfo verifyData = new DistroClientVerifyInfo(client.getClientId(),
                    client.getRevision());
            DistroKey distroKey = new DistroKey(client.getClientId(), TYPE);
            DistroData data = new DistroData(distroKey,
                    ApplicationUtils.getBean(Serializer.class).serialize(verifyData));
            data.setType(DataOperation.VERIFY);
            if (result == null) {
                result = new LinkedList<>();
            }
            result.add(data);
        }
    }
    return result;
}

方法中过滤了两个条件:client.isEphemeral() 确保仅处理临时客户端(持久客户端由 JRaft 协议管理),clientManager.isResponsibleClient(client) 确保仅收集本节点有责任管理的客户端(根据 DistroDataStorageImpl 的一致性哈希算法分配)。DistroClientVerifyInfo 的构建极为轻量——只取 clientIdrevision,序列化为字节数组后封装成 DistroData


四、processVerifyData:验证数据的接收处理

校验数据发送到目标节点后,接收端的验证链路是决定校验成败的关键。本节的时序图展示了从 DistroProtocol.onVerify()verifyClient() 的完整处理流程:

在这里插入图片描述

4.1 服务端接收入口

目标节点收到 DistroDataRequest 后,DistroDataRequestHandler 处理请求,根据请求头中的操作类型 DataOperation.VERIFY 路由到 DistroProtocol.onVerify()

com.alibaba.nacos.core.distributed.distro.DistroProtocol#onVerify

public boolean onVerify(DistroData distroData, String sourceAddress) {
    if (Loggers.DISTRO.isDebugEnabled()) {
        Loggers.DISTRO.debug("[DISTRO] Receive verify data type: {}, key: {}", distroData.getType(),
                distroData.getDistroKey());
    }
    String resourceType = distroData.getDistroKey().getResourceType();
    DistroDataProcessor dataProcessor = distroComponentHolder.findDataProcessor(resourceType);
    if (null == dataProcessor) {
        Loggers.DISTRO.warn("[DISTRO] Can't find verify data process for received data {}", resourceType);
        return false;
    }
    return dataProcessor.processVerifyData(distroData, sourceAddress);
}

onVerify()onReceive()(处理主动同步数据)共享 findDataProcessor() 查找逻辑,但其处理流程完全不同——主动同步走 processData(),校验走 processVerifyData()sourceAddress 参数由 DistroDataRequestHandler 从请求的 Source-Address 头中提取,在校验失败时用于日志记录(帮助运维人员定位哪个节点发来了不一致的校验数据)。

4.2 processVerifyData:反序列化与路由

com.alibaba.nacos.naming.consistency.ephemeral.distro.v2.DistroClientDataProcessor#processVerifyData

@Override
public boolean processVerifyData(DistroData distroData, String sourceAddress) {
    DistroClientVerifyInfo verifyData = ApplicationUtils.getBean(Serializer.class)
            .deserialize(distroData.getContent(), DistroClientVerifyInfo.class);
    if (clientManager.verifyClient(verifyData)) {
        return true;
    }
    Loggers.DISTRO.info("client {} is invalid, get new client from {}", verifyData.getClientId(), sourceAddress);
    return false;
}

方法的实现极为简洁:反序列化 → 委托 → 返回结果。反序列化将字节数组转换为 DistroClientVerifyInfo 对象(仅含 clientIdrevision),然后委托给 ClientManager 接口的 verifyClient() 方法。委托的设计使得不同的 ClientManager 实现可以定义不同的校验语义——同一个接口,ConnectionBasedClientManagerEphemeralIpPortClientManager 有各自的实现。

4.3 verifyClient 的两条实现策略

verifyClient() 有两个具体的实现,分别对应两种客户端管理模式:

ConnectionBasedClientManager(包路径 com.alibaba.nacos.naming.core.v2.client.manager.impl)适用于 gRPC 连接客户端:

com.alibaba.nacos.naming.core.v2.client.manager.impl.ConnectionBasedClientManager#verifyClient

@Override
public boolean verifyClient(DistroClientVerifyInfo verifyData) {
    ConnectionBasedClient client = clients.get(verifyData.getClientId());
    if (null != client) {
        if (0 == verifyData.getRevision() || client.getRevision() == verifyData.getRevision()) {
            client.setLastRenewTime();
            return true;
        } else {
            Loggers.DISTRO.info("[DISTRO-VERIFY-FAILED] ConnectionBasedClient[{}] revision local={}, remote={}",
                    client.getClientId(), client.getRevision(), verifyData.getRevision());
        }
    }
    return false;
}

EphemeralIpPortClientManager(适用于 HTTP/IP 端口客户端:

com.alibaba.nacos.naming.core.v2.client.manager.impl.EphemeralIpPortClientManager#verifyClient

@Override
public boolean verifyClient(DistroClientVerifyInfo verifyData) {
    String clientId = verifyData.getClientId();
    IpPortBasedClient client = clients.get(clientId);
    if (null != client) {
        if (0 == verifyData.getRevision() || client.getRevision() == verifyData.getRevision()) {
            NamingExecuteTaskDispatcher.getInstance()
                    .dispatchAndExecuteTask(clientId, new ClientBeatUpdateTask(client));
            return true;
        } else {
            Loggers.DISTRO.info("[DISTRO-VERIFY-FAILED] IpPortBasedClient[{}] revision local={}, remote={}",
                    client.getClientId(), client.getRevision(), verifyData.getRevision());
        }
    }
    return false;
}

两个实现的差异反映了客户端类型的业务特性差异,但其核心逻辑是一致的——都是一个"两段式判等"算法:

第一个判等点是 0 == verifyData.getRevision()向后兼容)。当校验请求来自旧版本节点时,revision 字段的值为 0(旧版本未实现 revision 机制)。此时无条件通过验证并刷新心跳时间。这个设计的精妙之处在于:版本升级时,旧节点和 2.4.3 新节点可以混布混跑,校验机制不会因为旧节点数据中缺少 revision 字段就错误地判定不一致。

第二个判等点是 client.getRevision() == verifyData.getRevision()精确匹配)。当双方都有 revision 时,比较本地和远端的版本号是否一致。一致表示数据一致,通过验证并刷新心跳时间;不一致表示数据有差异,校验失败记录日志。

两个实现的不同之处在于"通过验证后的行为":ConnectionBasedClientManager 调用 client.setLastRenewTime() 更新心跳时间;EphemeralIpPortClientManager 通过 NamingExecuteTaskDispatcher.dispatchAndExecuteTask() 提交 ClientBeatUpdateTask 来更新客户端心跳——后者使用了独立的调度器来处理心跳更新,与 gRPC 客户端的直连更新形成对比。

五、校验失败的冲突修复机制

processVerifyData() 返回 false(即 verifyClient() 失败)时,校验发起方需要执行补偿同步来修复数据不一致。本节的时序图展示了从接收到校验失败的响应到完成补偿同步的完整修复链路:

在这里插入图片描述

5.1 DistroVerifyCallbackWrapper:传输层回调的分层处理

校验请求发送后的异步回调由 DistroVerifyCallbackWrapper 处理,它是 RequestCallBack<Response> 的实现类。其 onResponse() 方法的处理逻辑分为两条路径:

com.alibaba.nacos.naming.consistency.ephemeral.distro.v2.DistroClientTransportAgent.DistroVerifyCallbackWrapper#onResponse

@Override
public void onResponse(Response response) {
    if (checkResponse(response)) {
        NamingTpsMonitor.distroVerifySuccess(member.getAddress(), member.getIp());
        distroCallback.onSuccess();
    } else {
        Loggers.DISTRO.info("Target {} verify client {} failed, sync new client", targetServer, clientId);
        NotifyCenter.publishEvent(new ClientEvent.ClientVerifyFailedEvent(clientId, targetServer));
        NamingTpsMonitor.distroVerifyFail(member.getAddress(), member.getIp());
        distroCallback.onFailed(null);
    }
}

成功路径:响应码为 SUCCESS 时,记录 TPS 成功指标后回调 distroCallback.onSuccess()DistroVerifyCallback.onSuccess() 仅输出 debug 日志,不做任何额外操作——校验通过是正常状态,不需要任何补偿动作。

失败路径:响应码非 SUCCESS 时(即接收端的 processVerifyData() 返回了 false),执行三个步骤——发布 ClientVerifyFailedEvent、记录 TPS 失败指标、回调 distroCallback.onFailed(null)。其中发布事件是最关键的一步,它是整个冲突修复链路的起点。

5.2 ClientVerifyFailedEvent:校验失败的事件载体

ClientVerifyFailedEvent 是连接校验失败与补偿同步的桥梁:

com.alibaba.nacos.naming.core.v2.event.client.ClientEvent.ClientVerifyFailedEvent

public static class ClientVerifyFailedEvent extends ClientEvent {
    
    private static final long serialVersionUID = 2023951686223780851L;
    private final String clientId;
    private final String targetServer;
    
    public ClientVerifyFailedEvent(String clientId, String targetServer) {
        super(null);
        this.clientId = clientId;
        this.targetServer = targetServer;
    }
    
    public String getClientId() {
        return clientId;
    }
    
    public String getTargetServer() {
        return targetServer;
    }
}

事件携带两个关键字段:clientId(本节点上校验发现不一致的客户端 ID)和 targetServer(目标节点地址,即刚刚校验失败的那个节点的地址)。NotifyCenter.publishEvent() 发布事件后,所有注册了 ClientEvent.ClientVerifyFailedEvent 类型监听的订阅者都会收到通知。

5.3 syncToVerifyFailedServer:补偿同步的执行

DistroClientDataProcessor 同时实现了 SmartSubscriberDistroDataProcessor,其 subscribeTypes() 注册了三种事件类型,其中就包括 ClientVerifyFailedEvent

com.alibaba.nacos.naming.consistency.ephemeral.distro.v2.DistroClientDataProcessor#subscribeTypes

@Override
public List<Class<? extends Event>> subscribeTypes() {
    List<Class<? extends Event>> result = new LinkedList<>();
    result.add(ClientEvent.ClientChangedEvent.class);
    result.add(ClientEvent.ClientDisconnectEvent.class);
    result.add(ClientEvent.ClientVerifyFailedEvent.class);
    return result;
}

onEvent() 方法根据事件类型分流:ClientChangedEventClientDisconnectEventsyncToAllServer(),而 ClientVerifyFailedEventsyncToVerifyFailedServer()

com.alibaba.nacos.naming.consistency.ephemeral.distro.v2.DistroClientDataProcessor#onEvent

@Override
public void onEvent(Event event) {
    if (EnvUtil.getStandaloneMode()) {
        return;
    }
    if (event instanceof ClientEvent.ClientVerifyFailedEvent) {
        syncToVerifyFailedServer((ClientEvent.ClientVerifyFailedEvent) event);
    } else {
        syncToAllServer((ClientEvent) event);
    }
}

private void syncToVerifyFailedServer(ClientEvent.ClientVerifyFailedEvent event) {
    Client client = clientManager.getClient(event.getClientId());
    if (isInvalidClient(client)) {
        return;
    }
    DistroKey distroKey = new DistroKey(client.getClientId(), TYPE);
    distroProtocol.syncToTarget(distroKey, DataOperation.ADD, event.getTargetServer(), 0L);
}

补偿同步的执行逻辑是:获取本地的客户端实例,检查是否合法(非空、临时、有责任管理),然后调用 distroProtocol.syncToTarget() 将客户端的全量数据同步到校验失败的目标节点。DataOperation.ADD 表示"新增/全量替换",0L 表示零延迟——不走延迟窗口,立即执行。syncToTarget()sync() 的区别在于,后者会遍历所有非自身节点,而 syncToTarget() 仅发送到指定的目标节点——这保证了补偿同步只修复"确实有差异"的节点,不会像 sync() 那样向所有节点广播。

补偿同步设置为零延迟的设计意图是明确的:校验发现不一致说明数据已经存在差异,窗口等待只会让不一致持续更久。立即同步虽然可能触发了多轮校验周期各修复一次(多个校验周期均发现不一致,各触发一次补偿同步),但增量同步算法的"比对差异"保证了重复同步不会产生副作用——如果数据已经一致,upgradeClient()equals() 比对会跳过更新。

5.4 DistroVerifyCallback:轻量级的日志回调

DistroVerifyExecuteTask 内部定义的 DistroVerifyCallback 是上层回调:

com.alibaba.nacos.core.distributed.distro.task.verify.DistroVerifyExecuteTask.DistroVerifyCallback

private class DistroVerifyCallback implements DistroCallback {
    
    @Override
    public void onSuccess() {
        if (Loggers.DISTRO.isDebugEnabled()) {
            Loggers.DISTRO.debug("[DISTRO] verify data for type {} to {} success", resourceType, targetServer);
        }
    }
    
    @Override
    public void onFailed(Throwable throwable) {
        DistroRecord distroRecord = DistroRecordsHolder.getInstance().getRecord(resourceType);
        distroRecord.verifyFail();
        if (Loggers.DISTRO.isDebugEnabled()) {
            Loggers.DISTRO.debug("[DISTRO-FAILED] verify data for type {} to {} failed.", resourceType, targetServer, throwable);
        }
    }
}

onSuccess() 仅输出 debug 级别的日志,无其他额外操作——校验成功是正常状态。onFailed() 更新 DistroRecordsHolder 中的校验失败计数,供运维监控使用。注意:这里的校验失败指的是网络传输层面的问题(如超时、节点不可达),与上一节说的"revision 不一致"不同——revision 不一致由 DistroVerifyCallbackWrapper.onResponse() 判断响应码后处理,而这里处理的是传输层异常回调。

六、整体链路串联

6.1 校验与修复的完整闭环

将前文五个章节串联起来,一条完整的校验与修复闭环如下:

“定时调度(DistroVerifyTimedTask 每 5 秒)→ 收集校验数据(getVerifyData:clientId + revision)→ 逐节点发送(DistroVerifyExecuteTask)→ 接收端验证(DistroProtocol.onVerifyprocessVerifyDataverifyClient)→ revision 一致通过 / 不一致失败 → 失败事件(ClientVerifyFailedEvent)→ 补偿同步(syncToVerifyFailedServersyncToTarget,0 延迟)→ 目标节点全量数据更新”

这条闭环的核心参与类共 8 个,跨越四个模块(core 定时任务、core 传输代理、naming 数据存储、naming 客户端管理),构成了一个从"定期检查"到"发现问题到修复问题"的完整治理链路。

下面的示意图以顺时针闭环的方式展示了这 8 个核心类在 5 秒周期内的完整协作关系:

在这里插入图片描述

6.2 五层架构与职责分离

从功能分层的角度来看,校验与修复机制可以分为五个层次:

层次 核心类 职责
调度层 DistroVerifyTimedTask 以 5 秒为周期轮询数据存储,收集校验数据,分发校验任务
执行层 DistroVerifyExecuteTask + DistroVerifyCallback 逐条发送校验数据,记录校验失败计数
传输层 DistroClientTransportAgent + DistroVerifyCallbackWrapper 前置检查 → gRPC 发送 → 响应校验 → TPS 监控 → 失败事件发布
验证层 DistroProtocol.onVerifyDistroClientDataProcessor.processVerifyDataClientManager.verifyClient 反序列化 → revision 比较 → 两段式判等(向后兼容 + 精确匹配)
补偿层 ClientVerifyFailedEventsyncToVerifyFailedServerDistroProtocol.syncToTarget 事件驱动 → 全量补偿同步到校验失败节点

五层架构的设计核心在于"轻量校验 + 重量修复"——校验阶段只传输最小的数据量(clientId + revision),最小化网络开销;一旦发现不一致,补偿阶段发送全量数据,一步到位修复。校验的轻量确保了高频率(5 秒)不会对集群造成压力,修复的重量确保了"发现问题即彻底解决"。

下面的分层示意图展示了五层架构的职责划分——前四层完成轻量校验,第五层执行重量修复:

在这里插入图片描述

6.3 设计思想凝练

回顾校验与修复机制的完整设计,可以凝练出三个核心设计思想:

一是"定时校验 + revision 比对"。校验协议的核心是极简的 DistroClientVerifyInfo(仅 clientId + revision),无需携带全量负载。5 秒的轮询周期在及时性与开销之间取得了平衡——太短增加网络负载,太长则数据不一致的窗口期过大。

二是"失败事件驱动"。校验机制不采用"失败后立即重试验证"的策略,而是通过 ClientVerifyFailedEvent 事件驱动补偿同步。这个设计的精妙之处在于——验证失败说明数据确实不一致,重试验证只会得到相同的结果;而补偿同步直接发送全量数据,从根本上解决差异。验证与修复的解耦使得可以独立调整校验频率和修复策略。

三是"向后兼容"。revision=0 的判等策略保证了旧版本节点可以无缝参与校验周期。这看似简单的策略取舍,实际上涉及了一个重要的设计决策——是"强制所有节点升级到支持 revision 的新版本"还是"兼容旧版本"。Nacos 选择了后者,保证了集群在滚动升级过程中校验机制不会因为版本混部而产生误报。

全文小结

本文聚焦 Distro 协议的容错处理机制,从定时校验调度、校验数据收集、revision 比对算法到校验失败的冲突修复,完整分析了定时校验如何兜底保证 AP 模式的最终一致性。

在调度层,DistroVerifyTimedTask 以 5 秒为周期轮询数据存储,通过三道检查(初始化状态、校验数据非空、传输代理存在性)确保校验任务仅在安全的条件下执行。在执行和传输层,DistroVerifyExecuteTask 逐条发送轻量级 DistroClientVerifyInfo(仅 clientId + revision),DistroVerifyCallbackWrapper 在校验失败时负责发布 ClientVerifyFailedEvent。在验证层,processVerifyData 将校验数据路由到对应 ClientManagerverifyClient 方法,两段式判等算法(revision=0 向后兼容 + 精确匹配)完成数据一致性判定。在补偿层,ClientVerifyFailedEvent 触发 syncToVerifyFailedServer 以 0 延迟直接调用 syncToTarget 发送全量数据,一步到位修复数据差异。这套"定时校验 + revision 比对 + 失败事件驱动补偿同步"的机制,与主动异步同步完美互补,构成了 AP 模式下最终一致性的双重保障。


原创不易,如果本文对您有帮助,带来了些许灵感或启发,烦请动动小手点赞、关注、转发、收藏。这是作者持续更新的动力源泉,衷心感谢您的支持。我会尽量在工作之余,为大家带来更高品质的内容,努力保持周更。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐