Clipping 微信公众号

拆解 Nacos 3.x Distro 协议:一个 1000ms 延迟背后的分布式设计

by Fox爱分享 原文 ↗
Created: 2026-07-03

公众号名称:Fox爱分享

作者名称:Fox爱分享

发布时间:2026-07-03 07:00

上篇我们聊了注册中心的 6 个核心设计,有读者问:“Distro 到底是怎么把数据从一个节点同步到另一个节点的?1000ms 的延迟窗口是干嘛的?” 这篇文章,我们把 Distro 源码翻个底朝天。


开篇:为什么 Nacos 自己造了个 Distro?

做分布式一致性,市面上有现成的方案:Gossip、Raft、Paxos。Nacos 选了 CP,用了 JRaft(基于 Raft);AP 这边,本来也可以接 Gossip,但阿里团队决定自己设计 Distro。

原因很简单:Gossip 是”每个人都跟每个人聊”,Distro 是”每个人只负责自己那摊子事,然后告诉别人”。

在服务注册场景下,一个 Client 的实例变更不需要全网广播——只需要让”负责这个 Client”的节点把变更同步给其他节点就行。这个”责任分片”的设计,让 Distro 比 Gossip 的消息量少一个数量级。


一、整体架构:三层模型

Distro 的代码分布在 corenaming 两个模块:

┌──────────────────────────────────────────────────┐
│  DistroProtocol (core)         总调度器            │
│  sync() / onReceive() / onVerify() / onQuery()    │
├───────────────┬────────────────┬─────────────────┤
│ DistroTransportAgent            DistroDataStorage │
│ 发送 sync/verify/query         获取数据 / 快照     │
│ 到远程节点                       / 校验数据         │
├───────────────┴────────────────┴─────────────────┤
│  DistroClientDataProcessor (naming)               │
│  同时实现 DataStorage + DataProcessor             │
│  Type = "Nacos:Naming:v2:ClientData"              │
└──────────────────────────────────────────────────┘

来看 DistroProtocol 的核心方法签名,它暴露了 4 个入口给业务层调用:

// 1. 增量同步:延时合并后同步到所有其他节点
public void sync(DistroKey distroKey, DataOperation action) {
    sync(distroKey, action, DistroConfig.getInstance().getSyncDelayMillis());
}

// 2. 定向同步:直接同步到指定节点(校验失败时用)
public void syncToTarget(DistroKey distroKey, DataOperation action,
    String targetServer, long delay)

// 3. 接收数据:收到其他节点的同步数据,找对应的 Processor 处理
public boolean onReceive(DistroData distroData)

// 4. 接收校验:收到校验请求,对账本节点数据是否正确
public boolean onVerify(DistroData distroData, String sourceAddress)

DistroProtocol 本身不存数据、不发请求,它只是一个路由层:根据 resourceType 找到对应的 TransportAgentDataStorageDataProcessor,然后调度执行。

这个设计的好处是:Distro 协议和具体的业务数据解耦naming 模块的 Client 数据用 Distro 同步,理论上 config 模块也可以复用同一套协议框架,只需要注册自己的 Processor。


二、两阶段任务模型(这是最精妙的)

Distro 的同步不是”立即发送”,而是先排队,再发送。整个流程分两阶段:

阶段一:延时合并(DelayTask)

DistroProtocol.sync() 被调用后,并不会立刻发网络请求。它创建一个 DistroDelayTask,丢进 DistroDelayTaskExecuteEngine

public void sync(DistroKey distroKey, DataOperation action, long delay) {
    for (Member each : memberManager.allMembersWithoutSelf()) {
        syncToTarget(distroKey, action, each.getAddress(), delay);
    }
}

public void syncToTarget(DistroKey distroKey, DataOperation action,
    String targetServer, long delay) {
    DistroDelayTask distroDelayTask =
        new DistroDelayTask(distroKeyWithTarget, action, delay);
    distroTaskEngineHolder.getDelayTaskExecuteEngine()
        .addTask(distroKeyWithTarget, distroDelayTask);
}

这里有一个非常关键的设计:merge 机制

// DistroDelayTask.merge()
@Override
public void merge(AbstractDelayTask task) {
    DistroDelayTask oldTask = (DistroDelayTask) task;
    if (!action.equals(oldTask.getAction()) && createTime < oldTask.getCreateTime()) {
        action = oldTask.getAction();
        createTime = oldTask.getCreateTime();
    }
    setLastProcessTime(oldTask.getLastProcessTime());
}

当同一个 DistroKey 的多个变更在 1000ms 窗口内到达时,后续任务会合并到已存在的任务上,只保留最新的操作类型和创建时间。

这有什么意义?

假设一个 Client 在 100ms 内连续注册了 3 个实例,如果每次注册都发一次网络请求,就是 3 次。有了 merge,三次变更合并成一次,只发一个请求——合并的粒度是 DistroKey,也就是单个 Client

默认的延时窗口是 1000ms,可以通过 nacos.core.protocol.distro.data.sync.delayMs 调整。

阶段二:执行同步(ExecuteTask)

延时到期后,DistroDelayTaskProcessor 把 DelayTask 转成真正的执行任务:

public boolean process(NacosTask task) {
    DistroDelayTask distroDelayTask = (DistroDelayTask) task;
    switch (distroDelayTask.getAction()) {
        case DELETE:
            DistroSyncDeleteTask syncDeleteTask =
                new DistroSyncDeleteTask(distroKey, distroComponentHolder);
            distroTaskEngineHolder.getExecuteWorkersManager()
                .addTask(distroKey, syncDeleteTask);
            returntrue;
        case CHANGE:
        case ADD:
            DistroSyncChangeTask syncChangeTask =
                new DistroSyncChangeTask(distroKey, distroComponentHolder);
            distroTaskEngineHolder.getExecuteWorkersManager()
                .addTask(distroKey, syncChangeTask);
            returntrue;
    }
}

DistroSyncChangeTask 做三件事:

// 1. 从 DataStorage 取出当前 Client 的最新数据(序列化)
DistroData distroData = getDistroData(type);

// 2. 通过 TransportAgent 发送到目标节点
getDistroComponentHolder().findTransportAgent(type)
    .syncData(distroData, getDistroKey().getTargetServer());

TransportAgent 在 naming 模块的实现是 DistroClientTransportAgent,底层用的是 gRPC(ClusterRpcClientProxy),支持同步和异步两种发送模式,超时默认 3000ms

总结两阶段模型

数据变更 → DistroDelayTask(1000ms延时+合并) → DistroSyncChangeTask → gRPC发送

这个设计用”1000ms 的延迟”换来了”消息量的几何级减少”。在微服务频繁发布/缩容的场景下,效果非常明显。


三、增量同步:事件驱动 + 责任分片

Distro 的增量同步是靠事件触发的,不是轮询。

DistroClientDataProcessor 同时注册为 SmartSubscriber,监听三种事件:

@Override
public List> subscribeTypes() {
    result.add(ClientEvent.ClientChangedEvent.class);       // Client 数据变更
    result.add(ClientEvent.ClientDisconnectEvent.class);    // Client 断开连接
    result.add(ClientEvent.ClientVerifyFailedEvent.class);  // 校验发现数据不一致
    return result;
}

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

ClientChangedEvent 触发时:

if (event instanceof ClientEvent.ClientChangedEvent) {
    DistroKey distroKey = new DistroKey(client.getClientId(), TYPE);
    distroProtocol.sync(distroKey, DataOperation.CHANGE);
}

但这里有一个关键判断——isInvalidClient()

private boolean isInvalidClient(Client client) {
    return null == client || !client.isEphemeral()    // 只同步临时实例
        || !clientManager.isResponsibleClient(client); // 只管自己负责的
}

这就是 Distro 的”责任分片”机制:每个 Nacos 节点只负责一部分 Client(通过 clientManager.isResponsibleClient(client) 判断)。只有当自己是该 Client 的”责任人”时,才会把变更同步到其他节点。

这个设计避免了”A 改了 B 也改,互相打架”的问题——写请求只有一个来源,就是责任节点


四、定时校验:版本号对账 + 自动修复

增量同步覆盖正常流程,但万一网络丢包、节点短暂故障呢?Distro 用定时校验兜底。

DistroVerifyTimedTask 每隔 5000ms(默认)执行一次:

@Override
public void run() {
    List targetServer = serverMemberManager.allMembersWithoutSelf();
    for (String each : distroComponentHolder.getDataStorageTypes()) {
        verifyForDataStorage(each, targetServer);
    }
}

private void verifyForDataStorage(String type, List targetServer) {
    DistroDataStorage dataStorage = distroComponentHolder.findDataStorage(type);
    if (!dataStorage.isFinishInitial()) return;
    
    List verifyData = dataStorage.getVerifyData();
    for (Member member : targetServer) {
        executeTaskExecuteEngine.addTask(member.getAddress() + type,
            new DistroVerifyExecuteTask(agent, verifyData, member.getAddress(), type));
    }
}

校验数据长什么样?

// DistroClientVerifyInfo
public class DistroClientVerifyInfo {
    private String clientId;   // "192.168.1.1#8888#true"
    private long revision;     // 数据版本号
}

只传 clientId + revision不传全量数据。校验逻辑在 DistroClientDataProcessor.processVerifyData()

public boolean processVerifyData(DistroData distroData, String sourceAddress) {
    DistroClientVerifyInfo verifyData = deserialize(distroData.getContent());
    if (clientManager.verifyClient(verifyData)) {
        return true;   // 版本号一致,数据正确
    }
    // 版本号不一致!需要从源节点拉取最新数据
    return false;
}

如果校验失败,DistroVerifyCallback 会发布 ClientVerifyFailedEvent

// DistroClientTransportAgent.DistroVerifyCallback
public void onResponse(Response response) {
    if (checkResponse(response)) {
        distroCallback.onSuccess();
    } else {
        // 校验失败 → 发布事件 → onEvent 里会 syncToVerifyFailedServer
        NotifyCenter.publishEvent(
            new ClientEvent.ClientVerifyFailedEvent(clientId, targetServer));
        distroCallback.onFailed(null);
    }
}

syncToVerifyFailedServer 不经过 DelayTask 队列,直接 0 延迟发送

private void syncToVerifyFailedServer(ClientEvent.ClientVerifyFailedEvent event) {
    // Verify failed data should be sync directly.
    distroProtocol.syncToTarget(distroKey, DataOperation.ADD,
        event.getTargetServer(), 0L);  // delay = 0
}

校验链路总结

每 5s → 收集所有本节点负责的Client的(clientId+revision)
→ 发校验请求到其他节点 → 对端对比版本号
→ 一致:OK | 不一致:发 ClientVerifyFailedEvent
→ 本节点收到事件 → 直接全量同步到不一致的节点(delay=0)


五、新节点加入:全量快照加载

一个新 Nacos 节点启动,内存是空的,它怎么知道自己应该负责哪些 Client?

DistroLoadDataTask

public DistroProtocol(...) {
    this.memberManager = memberManager;
    // ...
    startDistroTask();
}

private void startDistroTask() {
    if (EnvUtil.getStandaloneMode()) {
        isInitialized = true;
        return;
    }
    startVerifyTask();   // 启动定时校验
    startLoadTask();     // 启动全量加载
}

加载过程:

public void run() {
    try {
        load();  // 从已有节点拉快照
        if (!checkCompleted()) {
            // 没全部加载成功,30s 后重试
            GlobalExecutor.submitLoadDataTask(this,
                distroConfig.getLoadDataRetryDelayMillis());
        } else {
            loadCallback.onSuccess();
            Loggers.DISTRO.info("[DISTRO-INIT] load snapshot data success");
        }
    } catch (Exception e) {
        loadCallback.onFailed(e);
    }
}

具体拉取逻辑:

private boolean loadAllDataSnapshotFromRemote(String resourceType) {
    for (Member each : memberManager.allMembersWithoutSelf()) {
        // 1. 向已有节点请求全量快照
        DistroData distroData = transportAgent.getDatumSnapshot(each.getAddress());
        // 2. 交给 Processor 处理快照数据
        boolean result = dataProcessor.processSnapshot(distroData);
        if (result) {
            // 3. 标记初始化完成
            distroComponentHolder.findDataStorage(resourceType).finishInitial();
            returntrue;
        }
    }
    returnfalse;
}

getDatumSnapshot() 在服务端触发的路径是:

DistroProtocol.onSnapshot() → DistroClientDataProcessor.getDatumSnapshot()
→ 遍历所有临时Client → client.generateSyncData()
→ 序列化 → ClientSyncDatumSnapshot

加载完成后,finishInitial()isFinishInitial 置为 true。之后定时校验任务才会开始向其他节点发送校验数据。

关键参数

参数默认值含义
loadDataRetryDelay30s全量加载失败后的重试间隔
loadDataTimeout30s单次全量加载请求的超时时间

六、失败重试:不是指数退避,是固定延迟

Distro 的失败重试机制非常简洁。DistroClientTaskFailedHandler

@Override
public void retry(DistroKey distroKey, DataOperation action) {
    DistroDelayTask retryTask = new DistroDelayTask(distroKey, action,
        DistroConfig.getInstance().getSyncRetryDelayMillis());
    distroTaskEngineHolder.getDelayTaskExecuteEngine()
        .addTask(distroKey, retryTask);
}

就是重新创建一个 DelayTask,丢回延时队列。重试间隔默认 3000ms,固定值,不做指数退避。

为什么不退避?

因为 Distro 的失败场景通常是”目标节点暂时不可达”,不是”系统过载”。节点恢复后应该尽快同步,退避反而延迟了数据一致的时间窗口。这个设计假设”故障是短暂的,快速重试比慢慢退避更合理”。


七、完整时序:一次注册的 Distro 全链路

把前面讲的串起来,当一个临时实例注册到 Nacos 集群时,Distro 做了什么:

T=0     Client gRPC 注册实例

T=0     EphemeralClientOperationServiceImpl.registerInstance()
          → Client.addServiceInstance()
          → NotifyCenter.publishEvent(ClientChangedEvent)

T=0     DistroClientDataProcessor.onEvent()
          → isResponsibleClient? YES
          → distroProtocol.sync(clientId, CHANGE)

T=0     DistroDelayTask 入队 (delay=1000ms)

0-1000  同一 Client 的后续变更 → merge 到已有 DelayTask

T=1000  DistroDelayTaskProcessor 处理
          → DistroSyncChangeTask 创建
          → getDistroData(): 序列化 Client 全量数据
          → TransportAgent.syncData(): gRPC 发送到其他节点

T=1xxx  对端 DistroDataRequestHandler.handle()
          → distroProtocol.onReceive()
          → DistroClientDataProcessor.processData()
          → clientManager.syncClientConnected() + upgradeClient()

每 5s   DistroVerifyTimedTask 执行
          → getVerifyData(): 收集(clientId + revision)
          → 发送 verify 请求到其他节点
          → 版本对账 → 不一致则触发全量重新同步

总结:Distro 的 5 个核心设计哲学

设计哲学具体实现效果
延时合并DelayTask 1000ms + merge()高频变更不冲击网络
责任分片isResponsibleClient()每个 Client 只有一个写入口
版本对账VerifyTimedTask 每 5s + revision兜底修复不一致
全量快照DistroLoadDataTask + retry 30s新节点平滑加入
协议解耦ComponentHolder + SPI 模式同框架支撑多业务类型

最后说一句

Distro 不是 Gossip,不是 Raft,不是某种”新型共识算法”。它是一个为服务注册场景深度定制的 AP 同步协议

它的核心思路可以浓缩成一句话:

每个人只负责自己那摊子事,然后告诉别人。

“延时合并”减少消息量,“责任分片”避免写冲突,“版本对账”兜底修复,“全量快照”平滑加入——这四个机制组合在一起,构成了 Nacos 3.x AP 模式的底座。

这个底座上承载了数万节点规模的阿里内部集群,是经过生产验证的硬核设计。


关注「Fox爱分享」,持续输出分布式/云原生/AI 工程硬核内容。

下篇预告:Nacos 3.x AI 注册中心——Prompt、MCP、A2A 的资源管理究竟怎么做的?


内容效果不满意?点此反馈

输入关键词开始搜索