ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

四、Nacos源码系列:Nacos服务注册流程(三)

2026/9/1 11:00:22 拓冰建站 浏览量
四、Nacos源码系列:Nacos服务注册流程(三)

目录

一、服务端处理服务变更事件

1.1、任务执行引擎

1.1.1、NacosTaskExecuteEngine

1.1.2、AbstractNacosTaskExecuteEngine

1.1.3、NacosDelayTaskExecuteEngine

1.1.4、PushDelayTaskExecuteEngine

1.2、执行任务分发器

二、执行推送任务


一、服务端处理服务变更事件

接着上一篇文章,我们继续查找服务变更事件ServiceChangedEvent的订阅者,发现NamingSubscriberServiceV2Impl这个类订阅了ServiceChangedEvent,查看其onEvent()方法:

public void onEvent(Event event) {/*** 1.服务变更事件,推送这个变更给所有的订阅者* 2.如果服务被某个客户端订阅了,那么只会推送这个变更给这个客户端*/if (event instanceof ServiceEvent.ServiceChangedEvent) {// 服务变更时,向所有客户端推送ServiceEvent.ServiceChangedEvent serviceChangedEvent = (ServiceEvent.ServiceChangedEvent) event;// 获取Service服务信息Service service = serviceChangedEvent.getService();// 将service信息包装成一个PushDelayTask,添加到延时任务调度引擎,最终实际会执行PushDelayTaskProcessor.process()方法delayTaskEngine.addTask(service, new PushDelayTask(service, PushConfig.getInstance().getPushTaskDelay()));MetricsMonitor.incrementServiceChangeCount(service.getNamespace(), service.getGroup(), service.getName());} else if (event instanceof ServiceEvent.ServiceSubscribedEvent) {ServiceEvent.ServiceSubscribedEvent subscribedEvent = (ServiceEvent.ServiceSubscribedEvent) event;Service service = subscribedEvent.getService();delayTaskEngine.addTask(service, new PushDelayTask(service, PushConfig.getInstance().getPushTaskDelay(),subscribedEvent.getClientId()));}
}

 最关键的逻辑就是在这个添加任务方法:

delayTaskEngine.addTask(service, new PushDelayTask(service, PushConfig.getInstance().getPushTaskDelay()));

我们看下NamingSubscriberServiceV2Impl类的重要属性和构造方法:

// 任务引擎,在构造方法中进行初始化,类型为PushDelayTaskExecuteEngine
private final PushDelayTaskExecuteEngine delayTaskEngine;public NamingSubscriberServiceV2Impl(ClientManagerDelegate clientManager,ClientServiceIndexesManager indexesManager, ServiceStorage serviceStorage,NamingMetadataManager metadataManager, PushExecutorDelegate pushExecutor, SwitchDomain switchDomain) {this.clientManager = clientManager;this.indexesManager = indexesManager;// 初始化任务引擎this.delayTaskEngine = new PushDelayTaskExecuteEngine(clientManager, indexesManager, serviceStorage,metadataManager, pushExecutor, switchDomain);// 往NotifyCenter中注册自己NotifyCenter.registerSubscriber(this, NamingEventPublisherFactory.getInstance());}

首先我们介绍一下Nacos的任务引擎的设计。

1.1、任务执行引擎

先看下PushDelayTaskExecuteEngine的类图,如下所示。它实现了NacosTaskEcecuteEngine接口,继承自NacosDelayTaskExecuteEngine,而NacosDelayTaskExecuteEngine又继承自AbstractNacosTaskExecuteEngine。

1.1.1、NacosTaskExecuteEngine

先看下接口NacosTaskExecuteEngine。它定义了对NacosTask和NacosTaskProcessor的操作。 

public interface NacosTaskExecuteEngine<T extends NacosTask> extends Closeable {/*** 获取任务大小** @return size of task*/int size();/*** 判断任务引擎是否没有任务执行** @return true if the execute engine has no task to do, otherwise false*/boolean isEmpty();/*** 往任务引擎中添加处理类** @param key           key of task* @param taskProcessor task processor*/void addProcessor(Object key, NacosTaskProcessor taskProcessor);/*** 从任务引擎中删除处理类** @param key key of task*/void removeProcessor(Object key);/*** 从任务引擎中找到合适的处理类,没有找到的话,将使用默认的处理类** @param key key of task* @return task processor for task key or default processor if task processor for task key non-exist*/NacosTaskProcessor getProcessor(Object key);/*** 获取所有的处理类key** @return collection of processors*/Collection<Object> getAllProcessorKey();/*** 设置默认的处理类* task.** @param defaultTaskProcessor default task processor*/void setDefaultTaskProcessor(NacosTaskProcessor defaultTaskProcessor);/*** 往引擎中添加任务** @param key  key of task* @param task task*/void addTask(Object key, T task);/*** 从引擎中删除任务** @param key key of task* @return nacos task*/T removeTask(Object key);/*** 获取所有的任务Key** @return collection of task keys.*/Collection<Object> getAllTaskKeys();
}

1.1.2、AbstractNacosTaskExecuteEngine

接下来我们看下抽象的任务引擎类AbstractNacosTaskExecuteEngine。它有两个私有变量,分别是ConcurrentHashMap类型的taskProcessors,这个是对处理类NacosTaskProcessor的缓存。另一个是NacosTaskProcessor,这是一个默认的处理类,如果处理类缓存中不存在的话,就用这个处理类去处理。

其源码如下:

public abstract class AbstractNacosTaskExecuteEngine<T extends NacosTask> implements NacosTaskExecuteEngine<T> {private final Logger log;/*** 对处理类NacosTaskProcessor的缓存* key: Service服务* value: NacosTaskProcessor处理类*/private final ConcurrentHashMap<Object, NacosTaskProcessor> taskProcessors = new ConcurrentHashMap<>();/*** 默认的处理类。缓存中不存在的话,就用这个类去处理*/private NacosTaskProcessor defaultTaskProcessor;public AbstractNacosTaskExecuteEngine(Logger logger) {this.log = null != logger ? logger : LoggerFactory.getLogger(AbstractNacosTaskExecuteEngine.class.getName());}@Overridepublic void addProcessor(Object key, NacosTaskProcessor taskProcessor) {// 添加处理类,没有的时候才添加taskProcessors.putIfAbsent(key, taskProcessor);}@Overridepublic void removeProcessor(Object key) {// 移除处理类taskProcessors.remove(key);}@Overridepublic NacosTaskProcessor getProcessor(Object key) {// 如果缓存中能找到对应的处理器类,则直接返回。否则使用默认的处理类return taskProcessors.containsKey(key) ? taskProcessors.get(key) : defaultTaskProcessor;}@Overridepublic Collection<Object> getAllProcessorKey() {return taskProcessors.keySet();}@Overridepublic void setDefaultTaskProcessor(NacosTaskProcessor defaultTaskProcessor) {this.defaultTaskProcessor = defaultTaskProcessor;}protected Logger getEngineLog() {return log;}
}

1.1.3、NacosDelayTaskExecuteEngine

接下来我们看Nacos延迟任务执行引擎类NacosDelayTaskExecuteEngine,看名字就知道它是处理延迟任务的。

public class NacosDelayTaskExecuteEngine extends AbstractNacosTaskExecuteEngine<AbstractDelayTask> {/*** 定时任务线程池,在构造方法中初始化*/private final ScheduledExecutorService processingExecutor;/*** 任务队列* key:对应的服务*/protected final ConcurrentHashMap<Object, AbstractDelayTask> tasks;protected final ReentrantLock lock = new ReentrantLock();public NacosDelayTaskExecuteEngine(String name) {this(name, null);}public NacosDelayTaskExecuteEngine(String name, Logger logger) {this(name, 32, logger, 100L);}public NacosDelayTaskExecuteEngine(String name, Logger logger, long processInterval) {this(name, 32, logger, processInterval);}public NacosDelayTaskExecuteEngine(String name, int initCapacity, Logger logger) {this(name, initCapacity, logger, 100L);}public NacosDelayTaskExecuteEngine(String name, int initCapacity, Logger logger, long processInterval) {super(logger);// 初始化任务队列tasks = new ConcurrentHashMap<>(initCapacity);// 创建定时任务的线程池processingExecutor = ExecutorFactory.newSingleScheduledExecutorService(new NameThreadFactory(name));// 在指定的初始延迟时间(100毫秒)后开始执行任务,并按固定的时间间隔周期性(100毫秒)地执行任务。// 默认延时100毫秒执行ProcessRunnable,然后每隔100毫秒周期性执行ProcessRunnableprocessingExecutor.scheduleWithFixedDelay(new ProcessRunnable(), processInterval, processInterval, TimeUnit.MILLISECONDS);}@Overridepublic void addTask(Object key, AbstractDelayTask newTask) {// 加锁防并发处理,key就是对应的服务lock.lock();try {// ConcurrentHashMap<Object, AbstractDelayTask> tasks = new ConcurrentHashMap<>(initCapacity);// 通过key判断是否已存在map中AbstractDelayTask existTask = tasks.get(key);if (null != existTask) {// 服务存在的话,则需要合并任务,其实就是合并多个任务,一起执行newTask.merge(existTask);}// 将任务放入到map中,等待处理tasks.put(key, newTask);} finally {lock.unlock();}}/*** process tasks in execute engine.*/protected void processTasks() {Collection<Object> keys = getAllTaskKeys();for (Object taskKey : keys) {// 从队列中移除这个任务AbstractDelayTask task = removeTask(taskKey);if (null == task) {continue;}// taskKey示例值: Service{namespace='public', group='DEFAULT_GROUP', name='discovery-provider', ephemeral=true, revision=0}// 找到处理类NacosTaskProcessor processor = getProcessor(taskKey);if (null == processor) {getEngineLog().error("processor not found for task, so discarded. " + task);continue;}try {// ReAdd task if process failedif (!processor.process(task)) {// 处理失败的话,重新入队(即重试)retryFailedTask(taskKey, task);}} catch (Throwable e) {getEngineLog().error("Nacos task execute error ", e);retryFailedTask(taskKey, task);}}}private void retryFailedTask(Object key, AbstractDelayTask task) {task.setLastProcessTime(System.currentTimeMillis());addTask(key, task);}/*** 任务处理类*/private class ProcessRunnable implements Runnable {@Overridepublic void run() {try {processTasks();} catch (Throwable e) {getEngineLog().error(e.toString(), e);}}}
}

注意,NacosDelayTaskExecuteEngine的构造方法中创建了一个定时执行的线程池,默认延时100毫秒执行ProcessRunnable去执行任务,然后每隔100毫秒周期性执行。NacosDelayTaskExecuteEngine的addTask()还会将任务进行合并处理,提高效率。

1.1.4、PushDelayTaskExecuteEngine

最后看下PushDelayTaskExecuteEngine。在类的构造方法中就设置了默认的任务处理类PushDelayTaskProcessor。说明如果没有对应的key的任务处理类设置的话,就用这个处理类处理所有任务。

public class PushDelayTaskExecuteEngine extends NacosDelayTaskExecuteEngine {private final ClientManager clientManager;private final ClientServiceIndexesManager indexesManager;private final ServiceStorage serviceStorage;private final NamingMetadataManager metadataManager;private final PushExecutor pushExecutor;private final SwitchDomain switchDomain;public PushDelayTaskExecuteEngine(ClientManager clientManager, ClientServiceIndexesManager indexesManager,ServiceStorage serviceStorage, NamingMetadataManager metadataManager,PushExecutor pushExecutor, SwitchDomain switchDomain) {super(PushDelayTaskExecuteEngine.class.getSimpleName(), Loggers.PUSH);this.clientManager = clientManager;this.indexesManager = indexesManager;this.serviceStorage = serviceStorage;this.metadataManager = metadataManager;this.pushExecutor = pushExecutor;this.switchDomain = switchDomain;// 设置默认的处理类setDefaultTaskProcessor(new PushDelayTaskProcessor(this));}@Overrideprotected void processTasks() {if (!switchDomain.isPushEnabled()) {return;}super.processTasks();}private static class PushDelayTaskProcessor implements NacosTaskProcessor {private final PushDelayTaskExecuteEngine executeEngine;public PushDelayTaskProcessor(PushDelayTaskExecuteEngine executeEngine) {this.executeEngine = executeEngine;}@Overridepublic boolean process(NacosTask task) {PushDelayTask pushDelayTask = (PushDelayTask) task;Service service = pushDelayTask.getService();// 添加执行任务,其实就是NacosExecuteTaskExecuteEngine,实际上添加进去的是一个线程,重点关注PushExecuteTask.run()方法NamingExecuteTaskDispatcher.getInstance().dispatchAndExecuteTask(service, new PushExecuteTask(service, executeEngine, pushDelayTask));return true;}}
}

在简单了解了上述几个类的大致功能后,我们再回来看文章开头的那句代码:

delayTaskEngine.addTask(service, new PushDelayTask(service, PushConfig.getInstance().getPushTaskDelay()));

往PushDelayTaskExecuteEngine引擎中添加任务:

// com.alibaba.nacos.common.task.engine.NacosDelayTaskExecuteEngine#addTask
public void addTask(Object key, AbstractDelayTask newTask) {// 加锁防并发处理,key就是对应的服务lock.lock();try {// ConcurrentHashMap<Object, AbstractDelayTask> tasks = new ConcurrentHashMap<>(initCapacity);// 通过key判断是否已存在map中AbstractDelayTask existTask = tasks.get(key);if (null != existTask) {// 服务存在的话,则需要合并任务,其实就是合并多个任务,一起执行newTask.merge(existTask);}// 将任务放入到map中,等待处理tasks.put(key, newTask);} finally {lock.unlock();}
}

到这里,简单梳理大致的流程:

  • 1、往PushDelayTaskExecuteEngine引擎中添加任务;
  • 2、NacosDelayTaskExecuteEngine的构造方法中启动了定时执行的线程池任务,每隔100毫秒执行一次,首次执行会延迟100毫秒;
  • 3、定时任务执行的方法是NacosDelayTaskExecuteEngine.ProcessRunnable#run()方法,其内部调用了processTasks()方法;
protected void processTasks() {Collection<Object> keys = getAllTaskKeys();for (Object taskKey : keys) {// 从队列中移除这个任务AbstractDelayTask task = removeTask(taskKey);if (null == task) {continue;}// taskKey示例值: Service{namespace='public', group='DEFAULT_GROUP', name='discovery-provider', ephemeral=true, revision=0}// 找到处理类NacosTaskProcessor processor = getProcessor(taskKey);if (null == processor) {getEngineLog().error("processor not found for task, so discarded. " + task);continue;}try {// ReAdd task if process failedif (!processor.process(task)) {// 处理失败的话,重新入队(即重试)retryFailedTask(taskKey, task);}} catch (Throwable e) {getEngineLog().error("Nacos task execute error ", e);retryFailedTask(taskKey, task);}}
}

遍历任务队列中的所有任务,获取对应的处理类,如果没找到对应的处理类,就用默认的处理类。最终处理方法会调用processor.process(task)。

  • 4、PushDelayTaskExecuteEngine类中构造方法中设置了默认的处理类为PushDelayTaskProcessor,所以前面的processor.process(task),其实就是执行PushDelayTaskProcessor的process()方法;
private static class PushDelayTaskProcessor implements NacosTaskProcessor {private final PushDelayTaskExecuteEngine executeEngine;public PushDelayTaskProcessor(PushDelayTaskExecuteEngine executeEngine) {this.executeEngine = executeEngine;}@Overridepublic boolean process(NacosTask task) {PushDelayTask pushDelayTask = (PushDelayTask) task;Service service = pushDelayTask.getService();// 添加执行任务,其实就是NacosExecuteTaskExecuteEngine,实际上添加进去的是一个线程,重点关注PushExecuteTask.run()方法NamingExecuteTaskDispatcher.getInstance().dispatchAndExecuteTask(service, new PushExecuteTask(service, executeEngine, pushDelayTask));return true;}
}

如上述代码,绕了一大圈,好像还没到真正的处理,这里又有一个添加任务的逻辑。其实作者这样设计的目的是为了任务和处理类的解耦,并且异步化的去执行,并且还能将延迟的任务合并一起处理,可以提高性能,吞吐量,减少IO和网络操作。

1.2、执行任务分发器

继续往下分析,在前面的处理方法中,核心逻辑就一句代码: 

NamingExecuteTaskDispatcher.getInstance().dispatchAndExecuteTask(service, new PushExecuteTask(service, executeEngine, pushDelayTask));

NamingExecuteTaskDispatcher,翻译过来,就是执行任务分发器。NamingExecuteTaskDispatcher.getInstance()使用到了单例模式,整个进程共享这个实例。在类内部存在一个NacosExecuteTaskExecuteEngine的成员变量。

public class NamingExecuteTaskDispatcher {// 饿汉式单例private static final NamingExecuteTaskDispatcher INSTANCE = new NamingExecuteTaskDispatcher();// NacosExecuteTaskExecuteEngine跟NacosDelayTaskExecuteEngine很类似,都是继承于AbstractNacosTaskExecuteEngine,只不过NacosDelayTaskExecuteEngine是带有延时功能// 立即执行的任务引擎private final NacosExecuteTaskExecuteEngine executeEngine;private NamingExecuteTaskDispatcher() {// 初始化任务执行引擎:NacosExecuteTaskExecuteEngineexecuteEngine = new NacosExecuteTaskExecuteEngine(EnvUtil.FUNCTION_MODE_NAMING, Loggers.SRV_LOG);}public static NamingExecuteTaskDispatcher getInstance() {return INSTANCE;}public void dispatchAndExecuteTask(Object dispatchTag, AbstractExecuteTask task) {// 往任务引擎添加任务executeEngine.addTask(dispatchTag, task);}public String workersStatus() {return executeEngine.workersStatus();}public void destroy() throws Exception {executeEngine.shutdown();}
}

NacosExecuteTaskExecuteEngine跟前面介绍的NacosDelayTaskExecuteEngine很类似,都是继承于AbstractNacosTaskExecuteEngine,只不过NacosDelayTaskExecuteEngine是带有延时功能,NacosExecuteTaskExecuteEngine是立即执行的任务引擎。

下面来看看NacosExecuteTaskExecuteEngine的代码:

public class NacosExecuteTaskExecuteEngine extends AbstractNacosTaskExecuteEngine<AbstractExecuteTask> {// 任务执行worker,在构造方法中进行创建和初始化private final TaskExecuteWorker[] executeWorkers;public NacosExecuteTaskExecuteEngine(String name, Logger logger) {this(name, logger, ThreadUtils.getSuitableThreadCount(1));}public NacosExecuteTaskExecuteEngine(String name, Logger logger, int dispatchWorkerCount) {super(logger);// worker创建和初始化executeWorkers = new TaskExecuteWorker[dispatchWorkerCount];for (int mod = 0; mod < dispatchWorkerCount; ++mod) {executeWorkers[mod] = new TaskExecuteWorker(name, mod, dispatchWorkerCount, getEngineLog());}}}

NamingExecuteTaskDispatcher.getInstance().dispatchAndExecuteTask()实际上最终调用的是NacosExecuteTaskExecuteEngine的addTask()方法:

public void dispatchAndExecuteTask(Object dispatchTag, AbstractExecuteTask task) {// 往任务引擎添加任务executeEngine.addTask(dispatchTag, task);
}// com.alibaba.nacos.common.task.engine.NacosExecuteTaskExecuteEngine#addTask
public void addTask(Object tag, AbstractExecuteTask task) {// 获取处理类NacosTaskProcessor processor = getProcessor(tag);if (null != processor) {// 不为空,就用对应的processor处理processor.process(task);return;}// 没有找到处理类的话, 就用公共的TaskExecuteWorker执行TaskExecuteWorker worker = getWorker(tag);worker.process(task);
}

在addTask()方法内部,首先获取处理类,如果找到了处理类就用对应的处理类进行处理;如果没找到,则用公共的TaskExecuteWorker去执行。

我们再看下TaskEcecuteWorker是个什么东西。

public final class TaskExecuteWorker implements NacosTaskProcessor, Closeable {/*** 队列最大大小为32768*/private static final int QUEUE_CAPACITY = 1 << 15;private final Logger log;private final String name;/*** 阻塞队列*/private final BlockingQueue<Runnable> queue;private final AtomicBoolean closed;private final InnerWorker realWorker;public TaskExecuteWorker(final String name, final int mod, final int total) {this(name, mod, total, null);}public TaskExecuteWorker(final String name, final int mod, final int total, final Logger logger) {this.name = name + "_" + mod + "%" + total;// 阻塞队列this.queue = new ArrayBlockingQueue<>(QUEUE_CAPACITY);this.closed = new AtomicBoolean(false);this.log = null == logger ? LoggerFactory.getLogger(TaskExecuteWorker.class) : logger;// 内部执行worker,实际上是一个线程realWorker = new InnerWorker(this.name);// 启动workerrealWorker.start();}public String getName() {return name;}@Overridepublic boolean process(NacosTask task) {if (task instanceof AbstractExecuteTask) {// 添加任务到阻塞队列中putTask((Runnable) task);}return true;}private void putTask(Runnable task) {try {queue.put(task);} catch (InterruptedException ire) {log.error(ire.toString(), ire);}}public int pendingTaskCount() {return queue.size();}/*** Worker status.*/public String status() {return name + ", pending tasks: " + pendingTaskCount();}@Overridepublic void shutdown() throws NacosException {queue.clear();closed.compareAndSet(false, true);realWorker.interrupt();}}

从源码可以看到,在TaskExecuteWorker构造方法启动了一个InnerWorker,InnerWorker其实是一个线程,必然有run()方法:

/*** Inner execute worker.*/
private class InnerWorker extends Thread {InnerWorker(String name) {setDaemon(false);setName(name);}@Overridepublic void run() {while (!closed.get()) {try {// 从阻塞队列获取任务,在process()方法中通过putTask()将任务存入到了阻塞队列中Runnable task = queue.take();long begin = System.currentTimeMillis();// 执行任务task.run();long duration = System.currentTimeMillis() - begin;if (duration > 1000L) {log.warn("task {} takes {}ms", task, duration);}} catch (Throwable e) {log.error("[TASK-FAILED] " + e, e);}}}
}

我们看到,又是从阻塞队列中获取任务来执行,没有任务的时候就阻塞在那里,好了,TaskExecuteWorker的功能也很清晰了。我们再回到NacosExecuteTaskExecuteEngine#addTask方法,之前提到,如果没有找到处理类的话,就会用公共的TaskExecuteWorker执行,worker.process(task)其实就是调用的TaskExecuteWorker#process方法,在process()方法中通过putTask()将任务存入到了阻塞队列中,这样阻塞队列就有任务了,然后TaskExecuteWorker.InnerWorker#run这个内部线程的run方法就能拿到任务出来执行。

从阻塞队列中直接拿到任务,就能直接执行task.run(),难道我们之前存入阻塞队列中的是一个线程么?我们来验证一下:

NamingExecuteTaskDispatcher.getInstance().dispatchAndExecuteTask(service, new PushExecuteTask(service, executeEngine, pushDelayTask));

添加的是PushExecuteTask类型的任务,查看其类图:

我们发现PushExecuteTask继承自AbstractExecuteTask,而AbstractExecuteTask刚好实现了Runnable,确定了我们往阻塞队列中存入的就是一个线程对象来的,所以能够直接拿出来就执行,其实从TaskExecuteWorker的阻塞队列这个成员变量的声明private final BlockingQueue<Runnable> queue也可以看出。

下面我们来看看,从阻塞队列中拿到任务后,都做了哪些事情。

也就是我们需要分析下PushExecuteTask#run()方法:

// com.alibaba.nacos.naming.push.v2.task.PushExecuteTask#run
public void run() {try {// 生成推送所需要的数据。包括服务信息、服务元数据信息PushDataWrapper wrapper = generatePushData();// 获取客户端管理类ClientManager clientManager = delayTaskEngine.getClientManager();// 获取所有客户端或指定的客户端for (String each : getTargetClientIds()) {// 获取每个客户端Client client = clientManager.getClient(each);if (null == client) {// 说明客户端已经断开连接continue;}// 获取到这个服务的所有订阅者Subscriber subscriber = client.getSubscriber(service);// skip if nullif (subscriber == null) {continue;}// 推送给客户端delayTaskEngine.getPushExecutor().doPushWithCallback(each, subscriber, wrapper,new ServicePushCallback(each, subscriber, wrapper.getOriginalData(), delayTask.isPushToAll()));}} catch (Exception e) {Loggers.PUSH.error("Push task for service" + service.getGroupedServiceName() + " execute failed ", e);// 失败重推delayTaskEngine.addTask(service, new PushDelayTask(service, 1000L));}
}

可以看到,首先生成推送所需要的数据,包括服务信息、服务元数据信息。然后获取所有的客户端,挨个遍历,拿到每个客户端的订阅者,调用delayTaskEngine.getPushExecutor().doPushWithCallback()推送最新的服务信息给它们。

二、执行推送任务

delayTaskEngine.getPushExecutor().doPushWithCallback(each, subscriber, wrapper,new ServicePushCallback(each, subscriber, wrapper.getOriginalData(), delayTask.isPushToAll()));

第一步是获取delayTaskEngine.getPushExecutor(),跟踪这个类分析,发现是由构造方法传入的。再往上跟踪,可以看到该类是归Spring托管的,在NamingSubscriberServiceV2Impl的构造方法中注入的,注入的是PushExecutorDelegate委托类,如下图:

public class PushExecutorDelegate implements PushExecutor {// rpc执行类,V2版本使用private final PushExecutorRpcImpl rpcPushExecuteService;// udp执行类,V1版本使用private final PushExecutorUdpImpl udpPushExecuteService;public PushExecutorDelegate(PushExecutorRpcImpl rpcPushExecuteService, PushExecutorUdpImpl udpPushExecuteService) {this.rpcPushExecuteService = rpcPushExecuteService;this.udpPushExecuteService = udpPushExecuteService;}@Overridepublic void doPush(String clientId, Subscriber subscriber, PushDataWrapper data) {getPushExecuteService(clientId, subscriber).doPush(clientId, subscriber, data);}@Overridepublic void doPushWithCallback(String clientId, Subscriber subscriber, PushDataWrapper data,NamingPushCallback callBack) {// 执行推送getPushExecuteService(clientId, subscriber).doPushWithCallback(clientId, subscriber, data, callBack);}private PushExecutor getPushExecuteService(String clientId, Subscriber subscriber) {Optional<SpiPushExecutor> result = SpiImplPushExecutorHolder.getInstance().findPushExecutorSpiImpl(clientId, subscriber);if (result.isPresent()) {return result.get();}// 根据连接的客户端id识别是由upd推送还是rpc推送// 判断客户端ID中是否包含"#"return clientId.contains(IpPortBasedClient.ID_DELIMITER) ? udpPushExecuteService : rpcPushExecuteService;}
}

PushExecutorDelegate是一个委托类, 根据连接的客户端id识别是由upd推送还是rpc推送,动态选择不同的处理类。因为我们分析的是Nacos2.0的源码,所以这里使用的rpcPushExecuteService方式。

查看PushExecutorRpcImpl#doPushWithCallback()方法:

public void doPushWithCallback(String clientId, Subscriber subscriber, PushDataWrapper data,NamingPushCallback callBack) {// 获取服务信息ServiceInfo actualServiceInfo = getServiceInfo(data, subscriber);callBack.setActualServiceInfo(actualServiceInfo);// 构建一个NotifySubscriberRequest,通过grpc向客户端发送信息pushService.pushWithCallback(clientId, NotifySubscriberRequest.buildNotifySubscriberRequest(actualServiceInfo),callBack, GlobalExecutor.getCallbackExecutor());
}public void pushWithCallback(String connectionId, ServerRequest request, PushCallBack requestCallBack,Executor executor) {// 拿到客户端的连接Connection connection = connectionManager.getConnection(connectionId);if (connection != null) {try {// 发送异步请求connection.asyncRequest(request, new AbstractRequestCallBack(requestCallBack.getTimeout()) {@Overridepublic Executor getExecutor() {return executor;}@Overridepublic void onResponse(Response response) {if (response.isSuccess()) {requestCallBack.onSuccess();} else {requestCallBack.onFail(new NacosException(response.getErrorCode(), response.getMessage()));}}@Overridepublic void onException(Throwable e) {requestCallBack.onFail(e);}});} catch (ConnectionAlreadyClosedException e) {connectionManager.unregister(connectionId);requestCallBack.onSuccess();} catch (Exception e) {Loggers.REMOTE_DIGEST.error("error to send push response to connectionId ={},push response={}", connectionId,request, e);requestCallBack.onFail(e);}} else {requestCallBack.onSuccess();}
}

这里构建了一个NotifySubscriberRequest。在前面的文章分析grpc的时候,我们说过,通过request的类型,可以找到对应的处理类。

通过类引用,我们找到了com.alibaba.nacos.client.naming.remote.gprc.NamingPushRequestHandler这个处理类来处理NotifySubscriberRequest这种请求。

public Response requestReply(Request request) {if (request instanceof NotifySubscriberRequest) {NotifySubscriberRequest notifyRequest = (NotifySubscriberRequest) request;// 处理服务信息serviceInfoHolder.processServiceInfo(notifyRequest.getServiceInfo());return new NotifySubscriberResponse();}return null;
}public ServiceInfo processServiceInfo(ServiceInfo serviceInfo) {String serviceKey = serviceInfo.getKey();if (serviceKey == null) {return null;}// 获取老的服务ServiceInfo oldService = serviceInfoMap.get(serviceInfo.getKey());if (isEmptyOrErrorPush(serviceInfo)) {//empty or error push, just ignorereturn oldService;}// 重新存入客户端缓存中serviceInfoMap.put(serviceInfo.getKey(), serviceInfo);// 对比下服务信息是否发生变更boolean changed = isChangedServiceInfo(oldService, serviceInfo);if (StringUtils.isBlank(serviceInfo.getJsonFromServer())) {serviceInfo.setJsonFromServer(JacksonUtils.toJson(serviceInfo));}MetricsMonitor.getServiceInfoMapSizeMonitor().set(serviceInfoMap.size());if (changed) {NAMING_LOGGER.info("current ips:({}) service: {} -> {}", serviceInfo.ipCount(), serviceInfo.getKey(),JacksonUtils.toJson(serviceInfo.getHosts()));// 如果发生改变,发送实例变更事件NotifyCenter.publishEvent(new InstancesChangeEvent(notifierEventScope, serviceInfo.getName(), serviceInfo.getGroupName(),serviceInfo.getClusters(), serviceInfo.getHosts()));// 磁盘缓存也写一份DiskCache.write(serviceInfo, cacheDir);}return serviceInfo;
}

到这里,Nacos服务注册的流程就算完整结束了,可以看到,整个过程还是很复杂的,需要我们多看几次,多调试一下,跟踪其执行流程。