拆解 Nacos 3.x Distro 协议:一个 1000ms 延迟背后的分布式设计
公众号名称: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 的代码分布在 core 和 naming 两个模块:
┌──────────────────────────────────────────────────┐
│ 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 找到对应的 TransportAgent、DataStorage、DataProcessor,然后调度执行。
这个设计的好处是: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。之后定时校验任务才会开始向其他节点发送校验数据。
关键参数:
| 参数 | 默认值 | 含义 |
|---|---|---|
loadDataRetryDelay | 30s | 全量加载失败后的重试间隔 |
loadDataTimeout | 30s | 单次全量加载请求的超时时间 |
六、失败重试:不是指数退避,是固定延迟
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 的资源管理究竟怎么做的?
内容效果不满意?点此反馈