• Nacos源码详解


    因为Eureka的闭源,Nacos成为了现在Spring Cloud微服务注册中心的主流方案,那咱们废话不多说,直接开始读源码(基于2.2.0版本)。

    首先我们知道Nacos是基于Springboot实现自动注入的,老规矩我们来到spring-cloud-alibaba-nacos-discovery包下面的META-INF\spring.factories配置文件:

    1. org.springframework.boot.autoconfigure.EnableAutoConfiguration=\
    2. com.alibaba.cloud.nacos.discovery.NacosDiscoveryAutoConfiguration,\
    3. com.alibaba.cloud.nacos.ribbon.RibbonNacosAutoConfiguration,\
    4. com.alibaba.cloud.nacos.endpoint.NacosDiscoveryEndpointAutoConfiguration,\
    5. com.alibaba.cloud.nacos.registry.NacosServiceRegistryAutoConfiguration,\
    6. com.alibaba.cloud.nacos.discovery.NacosDiscoveryClientConfiguration,\
    7. com.alibaba.cloud.nacos.discovery.reactive.NacosReactiveDiscoveryClientConfiguration,\
    8. com.alibaba.cloud.nacos.discovery.configclient.NacosConfigServerAutoConfiguration
    9. org.springframework.cloud.bootstrap.BootstrapConfiguration=\
    10. com.alibaba.cloud.nacos.discovery.configclient.NacosDiscoveryClientConfigServiceBootstrapConfiguration

    我们首先来到NacosDiscoveryAutoConfiguration

    1. @Configuration(proxyBeanMethods = false)
    2. @ConditionalOnDiscoveryEnabled
    3. @ConditionalOnNacosDiscoveryEnabled
    4. public class NacosDiscoveryAutoConfiguration {
    5. @Bean
    6. @ConditionalOnMissingBean
    7. public NacosDiscoveryProperties nacosProperties() {
    8. return new NacosDiscoveryProperties();
    9. }
    10. @Bean
    11. @ConditionalOnMissingBean
    12. public NacosServiceDiscovery nacosServiceDiscovery(
    13. NacosDiscoveryProperties discoveryProperties) {
    14. return new NacosServiceDiscovery(discoveryProperties);
    15. }
    16. }

    这个自动配置类注入了两个Bean,第一个Bean是从配置文件中读spring.cloud.nacos.discovery前缀的配置项,然后是第二个Bean

    1. public class NacosServiceDiscovery {
    2. private NacosDiscoveryProperties discoveryProperties;
    3. public NacosServiceDiscovery(NacosDiscoveryProperties discoveryProperties) {
    4. this.discoveryProperties = discoveryProperties;
    5. }
    6. /**
    7. * Return all instances for the given service.
    8. * @param serviceId id of service
    9. * @return list of instances
    10. * @throws NacosException nacosException
    11. */
    12. public List getInstances(String serviceId) throws NacosException {
    13. String group = discoveryProperties.getGroup();
    14. List instances = discoveryProperties.namingServiceInstance()
    15. .selectInstances(serviceId, group, true);
    16. return hostToServiceInstanceList(instances, serviceId);
    17. }
    18. /**
    19. * Return the names of all services.
    20. * @return list of service names
    21. * @throws NacosException nacosException
    22. */
    23. public List getServices() throws NacosException {
    24. String group = discoveryProperties.getGroup();
    25. ListView services = discoveryProperties.namingServiceInstance()
    26. .getServicesOfServer(1, Integer.MAX_VALUE, group);
    27. return services.getData();
    28. }
    29. public static List hostToServiceInstanceList(
    30. List instances, String serviceId) {
    31. List result = new ArrayList<>(instances.size());
    32. for (Instance instance : instances) {
    33. ServiceInstance serviceInstance = hostToServiceInstance(instance, serviceId);
    34. if (serviceInstance != null) {
    35. result.add(serviceInstance);
    36. }
    37. }
    38. return result;
    39. }
    40. public static ServiceInstance hostToServiceInstance(Instance instance,
    41. String serviceId) {
    42. if (instance == null || !instance.isEnabled() || !instance.isHealthy()) {
    43. return null;
    44. }
    45. NacosServiceInstance nacosServiceInstance = new NacosServiceInstance();
    46. nacosServiceInstance.setHost(instance.getIp());
    47. nacosServiceInstance.setPort(instance.getPort());
    48. nacosServiceInstance.setServiceId(serviceId);
    49. Map metadata = new HashMap<>();
    50. metadata.put("nacos.instanceId", instance.getInstanceId());
    51. metadata.put("nacos.weight", instance.getWeight() + "");
    52. metadata.put("nacos.healthy", instance.isHealthy() + "");
    53. metadata.put("nacos.cluster", instance.getClusterName() + "");
    54. metadata.putAll(instance.getMetadata());
    55. nacosServiceInstance.setMetadata(metadata);
    56. if (metadata.containsKey("secure")) {
    57. boolean secure = Boolean.parseBoolean(metadata.get("secure"));
    58. nacosServiceInstance.setSecure(secure);
    59. }
    60. return nacosServiceInstance;
    61. }

    根据代码和注释可以看出,这个Bean的核心功能是根据配置文件里Nacos的相关配置,获取实例和实例对象转化的一些功能,然后是NacosDiscoveryEndpointAutoConfiguration

    1. @Configuration(proxyBeanMethods = false)
    2. @ConditionalOnClass(Endpoint.class)
    3. @ConditionalOnNacosDiscoveryEnabled
    4. public class NacosDiscoveryEndpointAutoConfiguration {
    5. @Bean
    6. @ConditionalOnMissingBean
    7. @ConditionalOnEnabledEndpoint
    8. public NacosDiscoveryEndpoint nacosDiscoveryEndpoint(
    9. NacosDiscoveryProperties nacosDiscoveryProperties) {
    10. return new NacosDiscoveryEndpoint(nacosDiscoveryProperties);
    11. }
    12. @Bean
    13. @ConditionalOnEnabledHealthIndicator("nacos-discovery")
    14. public HealthIndicator nacosDiscoveryHealthIndicator(
    15. NacosDiscoveryProperties nacosDiscoveryProperties) {
    16. return new NacosDiscoveryHealthIndicator(
    17. nacosDiscoveryProperties.namingServiceInstance());
    18. }
    19. }

    这个类注入了两个Bean,第一个Bean NacosDiscoveryEndpoint 的作用是作为Nacos服务发现的端点,获取当前客户端的所有注册的服务。

    第二个Bean NacosDiscoveryHealthIndicator的作用是通过向Nacos注册中心请求/operator/metrics 接口来确认健康状态。

    接下来这个自动配置类NacosServiceRegistryAutoConfiguration比较核心

    1. @Configuration(proxyBeanMethods = false)
    2. @EnableConfigurationProperties
    3. /*是否开启Nacos服务发现*/
    4. @ConditionalOnNacosDiscoveryEnabled
    5. /*是否开启自动注册,默认为True*/
    6. @ConditionalOnProperty(value = "spring.cloud.service-registry.auto-registration.enabled",
    7. matchIfMissing = true)
    8. /*在以下几个配置类装配完成后才进行装配*/
    9. @AutoConfigureAfter({ AutoServiceRegistrationConfiguration.class,
    10. AutoServiceRegistrationAutoConfiguration.class,
    11. NacosDiscoveryAutoConfiguration.class })
    12. public class NacosServiceRegistryAutoConfiguration {
    13. /*读取配置文件*/
    14. @Bean
    15. public NacosServiceRegistry nacosServiceRegistry(
    16. NacosDiscoveryProperties nacosDiscoveryProperties) {
    17. return new NacosServiceRegistry(nacosDiscoveryProperties);
    18. }
    19. /*Nacos服务注册Bean*/
    20. @Bean
    21. @ConditionalOnBean(AutoServiceRegistrationProperties.class)
    22. public NacosRegistration nacosRegistration(
    23. NacosDiscoveryProperties nacosDiscoveryProperties,
    24. ApplicationContext context) {
    25. return new NacosRegistration(nacosDiscoveryProperties, context);
    26. }
    27. /*Nacos服务自动注册Bean*/
    28. @Bean
    29. @ConditionalOnBean(AutoServiceRegistrationProperties.class)
    30. public NacosAutoServiceRegistration nacosAutoServiceRegistration(
    31. NacosServiceRegistry registry,
    32. AutoServiceRegistrationProperties autoServiceRegistrationProperties,
    33. NacosRegistration registration) {
    34. return new NacosAutoServiceRegistration(registry,
    35. autoServiceRegistrationProperties, registration);
    36. }
    37. }

    该配置类装配了三个Bean,我们重点关注下第三个同来实现服务自动注册的Bean NacosAutoServiceRegistration

    1. public class NacosAutoServiceRegistration
    2. extends AbstractAutoServiceRegistration {
    3. private static final Logger log = LoggerFactory
    4. .getLogger(NacosAutoServiceRegistration.class);
    5. private NacosRegistration registration;
    6. public NacosAutoServiceRegistration(ServiceRegistry serviceRegistry,
    7. AutoServiceRegistrationProperties autoServiceRegistrationProperties,
    8. NacosRegistration registration) {
    9. super(serviceRegistry, autoServiceRegistrationProperties);
    10. this.registration = registration;
    11. }
    12. @Deprecated
    13. public void setPort(int port) {
    14. getPort().set(port);
    15. }
    16. @Override
    17. protected NacosRegistration getRegistration() {
    18. if (this.registration.getPort() < 0 && this.getPort().get() > 0) {
    19. this.registration.setPort(this.getPort().get());
    20. }
    21. Assert.isTrue(this.registration.getPort() > 0, "service.port has not been set");
    22. return this.registration;
    23. }
    24. @Override
    25. protected NacosRegistration getManagementRegistration() {
    26. return null;
    27. }
    28. @Override
    29. protected void register() {
    30. if (!this.registration.getNacosDiscoveryProperties().isRegisterEnabled()) {
    31. log.debug("Registration disabled.");
    32. return;
    33. }
    34. if (this.registration.getPort() < 0) {
    35. this.registration.setPort(getPort().get());
    36. }
    37. super.register();
    38. }
    39. @Override
    40. protected void registerManagement() {
    41. if (!this.registration.getNacosDiscoveryProperties().isRegisterEnabled()) {
    42. return;
    43. }
    44. super.registerManagement();
    45. }
    46. @Override
    47. protected Object getConfiguration() {
    48. return this.registration.getNacosDiscoveryProperties();
    49. }
    50. @Override
    51. protected boolean isEnabled() {
    52. return this.registration.getNacosDiscoveryProperties().isRegisterEnabled();
    53. }
    54. @Override
    55. @SuppressWarnings("deprecation")
    56. protected String getAppName() {
    57. String appName = registration.getNacosDiscoveryProperties().getService();
    58. return StringUtils.isEmpty(appName) ? super.getAppName() : appName;
    59. }

    通过名字我们可以看出,这个Bean的核心是redister方法,该方法调用的是父抽象类AbstractAutoServiceRegistration中的实现,然后在父类中它又调用了接口ServiceRegistry的register方法,实际由NacosServiceRegistry实现,我们进入到NacosServiceRegistry实现的register方法

    1. public void register(Registration registration) {
    2. if (StringUtils.isEmpty(registration.getServiceId())) {
    3. log.warn("No service to register for nacos client...");
    4. return;
    5. }
    6. String serviceId = registration.getServiceId();
    7. String group = nacosDiscoveryProperties.getGroup();
    8. Instance instance = getNacosInstanceFromRegistration(registration);
    9. try {
    10. namingService.registerInstance(serviceId, group, instance);
    11. log.info("nacos registry, {} {} {}:{} register finished", group, serviceId,
    12. instance.getIp(), instance.getPort());
    13. }
    14. catch (Exception e) {
    15. log.error("nacos registry, {} register failed...{},", serviceId,
    16. registration.toString(), e);
    17. // rethrow a RuntimeException if the registration is failed.
    18. // issue : https://github.com/alibaba/spring-cloud-alibaba/issues/1132
    19. rethrowRuntimeException(e);
    20. }
    21. }

    然后我们重点关注其中namingService.registerInstance方法,namingService是一个接口,由NacosNamingService实现,我们进入到方法实现

    1. public void registerInstance(String serviceName, String groupName, Instance instance) throws NacosException {
    2. //首先判断该实例是临时还是持久的,默认是临时的
    3. if (instance.isEphemeral()) {
    4. //新建心跳对象
    5. BeatInfo beatInfo = new BeatInfo();
    6. //设置serviceName
    7. beatInfo.setServiceName(NamingUtils.getGroupedName(serviceName, groupName));
    8. //设置IP
    9. beatInfo.setIp(instance.getIp());
    10. //设置端口号
    11. beatInfo.setPort(instance.getPort());
    12. //设置集群名
    13. beatInfo.setCluster(instance.getClusterName());
    14. //设置实例权重
    15. beatInfo.setWeight(instance.getWeight());
    16. beatInfo.setMetadata(instance.getMetadata());
    17. beatInfo.setScheduled(false);
    18. long instanceInterval = instance.getInstanceHeartBeatInterval();
    19. //设置发送心跳时间间隔,默认5秒
    20. beatInfo.setPeriod(instanceInterval == 0 ? DEFAULT_HEART_BEAT_INTERVAL : instanceInterval);
    21. //放入定时线程池等待执行
    22. beatReactor.addBeatInfo(NamingUtils.getGroupedName(serviceName, groupName), beatInfo);
    23. }
    24. //进行服务注册
    25. serverProxy.registerService(NamingUtils.getGroupedName(serviceName, groupName), groupName, instance);
    26. }

    该方法首先封装了一个心跳对象,并放入ScheduledThreadPoolExecutor线程池中,默认每5秒执行一次,然后调用注册中心代理serverProxy通过向注册中心/instance接口发送进行服务注册

    1. public void registerService(String serviceName, String groupName, Instance instance) throws NacosException {
    2. NAMING_LOGGER.info("[REGISTER-SERVICE] {} registering service {} with instance: {}",
    3. namespaceId, serviceName, instance);
    4. final Map params = new HashMap(9);
    5. params.put(CommonParams.NAMESPACE_ID, namespaceId);
    6. params.put(CommonParams.SERVICE_NAME, serviceName);
    7. params.put(CommonParams.GROUP_NAME, groupName);
    8. params.put(CommonParams.CLUSTER_NAME, instance.getClusterName());
    9. params.put("ip", instance.getIp());
    10. params.put("port", String.valueOf(instance.getPort()));
    11. params.put("weight", String.valueOf(instance.getWeight()));
    12. params.put("enable", String.valueOf(instance.isEnabled()));
    13. params.put("healthy", String.valueOf(instance.isHealthy()));
    14. params.put("ephemeral", String.valueOf(instance.isEphemeral()));
    15. params.put("metadata", JSON.toJSONString(instance.getMetadata()));
    16. reqAPI(UtilAndComs.NACOS_URL_INSTANCE, params, HttpMethod.POST);
    17. }

    我们来到注册中心的/instance接口,重点关注InstanceController的register方法

    1. @CanDistro
    2. @PostMapping
    3. @TpsControl(pointName = "NamingInstanceRegister", name = "HttpNamingInstanceRegister")
    4. @Secured(action = ActionTypes.WRITE)
    5. public String register(HttpServletRequest request) throws Exception {
    6. final String namespaceId = WebUtils.optional(request, CommonParams.NAMESPACE_ID,
    7. Constants.DEFAULT_NAMESPACE_ID);
    8. final String serviceName = WebUtils.required(request, CommonParams.SERVICE_NAME);
    9. NamingUtils.checkServiceNameFormat(serviceName);
    10. final Instance instance = HttpRequestInstanceBuilder.newBuilder()
    11. .setDefaultInstanceEphemeral(switchDomain.isDefaultInstanceEphemeral()).setRequest(request).build();
    12. getInstanceOperator().registerInstance(namespaceId, serviceName, instance);
    13. NotifyCenter.publishEvent(new RegisterInstanceTraceEvent(System.currentTimeMillis(), "", false, namespaceId,
    14. NamingUtils.getGroupName(serviceName), NamingUtils.getServiceName(serviceName), instance.getIp(),
    15. instance.getPort()));
    16. return "ok";
    17. }

    重点关注getInstanceOperator().registerInstance(namespaceId, serviceName, instance)方法

    1. /**
    2. * This method creates {@code IpPortBasedClient} if it doesn't exist.
    3. */
    4. @Override
    5. public void registerInstance(String namespaceId, String serviceName, Instance instance) throws NacosException {
    6. NamingUtils.checkInstanceIsLegal(instance);
    7. //判断是否临时
    8. boolean ephemeral = instance.isEphemeral();
    9. //获取客户端ID
    10. String clientId = IpPortBasedClient.getClientId(instance.toInetAddr(), ephemeral);
    11. //如果客户端不存在则新建
    12. createIpPortClientIfAbsent(clientId);
    13. //封装service对象
    14. Service service = getService(namespaceId, serviceName, ephemeral);
    15. //进行实例注册
    16. clientOperationService.registerInstance(service, instance, clientId);
    17. }

    clientOperationService为一个接口且有多个实现,我们查看其实现类找到EphemeralClientOperationServiceImpl中的registerInstance方法

    1. public void registerInstance(Service service, Instance instance, String clientId) throws NacosException {
    2. //检查实例是否合法
    3. NamingUtils.checkInstanceIsLegal(instance);
    4. //获取单例的service
    5. Service singleton = ServiceManager.getInstance().getSingleton(service);
    6. if (!singleton.isEphemeral()) {
    7. throw new NacosRuntimeException(NacosException.INVALID_PARAM,
    8. String.format("Current service %s is persistent service, can't register ephemeral instance.",
    9. singleton.getGroupedServiceName()));
    10. }
    11. //获取客户端
    12. Client client = clientManager.getClient(clientId);
    13. if (!clientIsLegal(client, clientId)) {
    14. return;
    15. }
    16. //封装实例发布信息
    17. InstancePublishInfo instanceInfo = getPublishInfo(instance);
    18. client.addServiceInstance(singleton, instanceInfo);
    19. client.setLastUpdatedTime();
    20. client.recalculateRevision();
    21. //通知中心发布客户端服务注册事件
    22. NotifyCenter.publishEvent(new ClientOperationEvent.ClientRegisterServiceEvent(singleton, clientId));
    23. NotifyCenter
    24. .publishEvent(new MetadataEvent.InstanceMetadataEvent(singleton, instanceInfo.getMetadataId(), false));
    25. }

    然后我们关注publishEvent

    1. /**
    2. * Request publisher publish event Publishers load lazily, calling publisher.
    3. *
    4. * @param eventType class Instances type of the event type.
    5. * @param event event instance.
    6. */
    7. private static boolean publishEvent(final Class eventType, final Event event) {
    8. if (ClassUtils.isAssignableFrom(SlowEvent.class, eventType)) {
    9. return INSTANCE.sharePublisher.publish(event);
    10. }
    11. final String topic = ClassUtils.getCanonicalName(eventType);
    12. EventPublisher publisher = INSTANCE.publisherMap.get(topic);
    13. if (publisher != null) {
    14. return publisher.publish(event);
    15. }
    16. if (event.isPluginEvent()) {
    17. return true;
    18. }
    19. LOGGER.warn("There are no [{}] publishers for this event, please register", topic);
    20. return false;
    21. }

    在DefaultPublisher中的publish方法,加入阻塞队列等待执行

    1. @Override
    2. public boolean publish(Event event) {
    3. checkIsStart();
    4. //加入阻塞队列尾部
    5. boolean success = this.queue.offer(event);
    6. if (!success) {
    7. LOGGER.warn("Unable to plug in due to interruption, synchronize sending time, event : {}", event);
    8. //接收事件
    9. receiveEvent(event);
    10. return true;
    11. }
    12. return true;
    13. }
    14. void checkIsStart() {
    15. if (!initialized) {
    16. throw new IllegalStateException("Publisher does not start");
    17. }
    18. }
    19. @Override
    20. public void shutdown() {
    21. this.shutdown = true;
    22. this.queue.clear();
    23. }
    24. public boolean isInitialized() {
    25. return initialized;
    26. }
    27. /**
    28. * Receive and notifySubscriber to process the event.
    29. *
    30. * @param event {@link Event}.
    31. */
    32. void receiveEvent(Event event) {
    33. final long currentEventSequence = event.sequence();
    34. //检查有无订阅者
    35. if (!hasSubscriber()) {
    36. LOGGER.warn("[NotifyCenter] the {} is lost, because there is no subscriber.", event);
    37. return;
    38. }
    39. // Notification single event listener
    40. for (Subscriber subscriber : subscribers) {
    41. if (!subscriber.scopeMatches(event)) {
    42. continue;
    43. }
    44. // Whether to ignore expiration events
    45. if (subscriber.ignoreExpireEvent() && lastEventSequence > currentEventSequence) {
    46. LOGGER.debug("[NotifyCenter] the {} is unacceptable to this subscriber, because had expire",
    47. event.getClass());
    48. continue;
    49. }
    50. // Because unifying smartSubscriber and subscriber, so here need to think of compatibility.
    51. // Remove original judge part of codes.
    52. //通知订阅者
    53. notifySubscriber(subscriber, event);
    54. }
    55. }
    56. @Override
    57. public void notifySubscriber(final Subscriber subscriber, final Event event) {
    58. LOGGER.debug("[NotifyCenter] the {} will received by {}", event, subscriber);
    59. final Runnable job = () -> subscriber.onEvent(event);
    60. final Executor executor = subscriber.executor();
    61. if (executor != null) {
    62. executor.execute(job);
    63. } else {
    64. try {
    65. //执行事件
    66. job.run();
    67. } catch (Throwable e) {
    68. LOGGER.error("Event callback exception: ", e);
    69. }
    70. }
    71. }

    然后我们通过ClientRegisterServiceEvent的其他引用找到了ClientServiceIndexesManager中的onEvent方法,发现ClientRegisterServiceEvent事件被封装成了ServiceChangedEvent并发布

    1. @Override
    2. public void onEvent(Event event) {
    3. if (event instanceof ClientOperationEvent.ClientReleaseEvent) {
    4. handleClientDisconnect((ClientOperationEvent.ClientReleaseEvent) event);
    5. } else if (event instanceof ClientOperationEvent) {
    6. handleClientOperation((ClientOperationEvent) event);
    7. }
    8. }
    9. private void handleClientDisconnect(ClientOperationEvent.ClientReleaseEvent event) {
    10. Client client = event.getClient();
    11. for (Service each : client.getAllSubscribeService()) {
    12. removeSubscriberIndexes(each, client.getClientId());
    13. }
    14. DeregisterInstanceReason reason = event.isNative()
    15. ? DeregisterInstanceReason.NATIVE_DISCONNECTED : DeregisterInstanceReason.SYNCED_DISCONNECTED;
    16. long currentTimeMillis = System.currentTimeMillis();
    17. for (Service each : client.getAllPublishedService()) {
    18. removePublisherIndexes(each, client.getClientId());
    19. InstancePublishInfo instance = client.getInstancePublishInfo(each);
    20. NotifyCenter.publishEvent(new DeregisterInstanceTraceEvent(currentTimeMillis,
    21. "", false, reason, each.getNamespace(), each.getGroup(), each.getName(),
    22. instance.getIp(), instance.getPort()));
    23. }
    24. }
    25. private void handleClientOperation(ClientOperationEvent event) {
    26. Service service = event.getService();
    27. String clientId = event.getClientId();
    28. if (event instanceof ClientOperationEvent.ClientRegisterServiceEvent) {
    29. addPublisherIndexes(service, clientId);
    30. } else if (event instanceof ClientOperationEvent.ClientDeregisterServiceEvent) {
    31. removePublisherIndexes(service, clientId);
    32. } else if (event instanceof ClientOperationEvent.ClientSubscribeServiceEvent) {
    33. addSubscriberIndexes(service, clientId);
    34. } else if (event instanceof ClientOperationEvent.ClientUnsubscribeServiceEvent) {
    35. removeSubscriberIndexes(service, clientId);
    36. }
    37. }
    38. private void addPublisherIndexes(Service service, String clientId) {
    39. publisherIndexes.computeIfAbsent(service, key -> new ConcurrentHashSet<>());
    40. publisherIndexes.get(service).add(clientId);
    41. NotifyCenter.publishEvent(new ServiceEvent.ServiceChangedEvent(service, true));
    42. }

    而Nacos注册中心是通过JRaftServer的NacosStateMachine中RequestProcessor对象去执行onApply方法,而RequestProcessor为抽象类,其实现类为InstanceMetadataProcessor,我们来到其中的onApply方法

    1. @Override
    2. public Response onApply(WriteRequest request) {
    3. readLock.lock();
    4. try {
    5. MetadataOperation op = serializer.deserialize(request.getData().toByteArray(), processType);
    6. switch (DataOperation.valueOf(request.getOperation())) {
    7. case ADD:
    8. case CHANGE:
    9. updateInstanceMetadata(op);
    10. break;
    11. case DELETE:
    12. deleteInstanceMetadata(op);
    13. break;
    14. default:
    15. return Response.newBuilder().setSuccess(false)
    16. .setErrMsg("Unsupported operation " + request.getOperation()).build();
    17. }
    18. return Response.newBuilder().setSuccess(true).build();
    19. } catch (Exception e) {
    20. Loggers.RAFT.error("onApply {} instance metadata operation failed. ", request.getOperation(), e);
    21. String errorMessage = null == e.getMessage() ? e.getClass().getName() : e.getMessage();
    22. return Response.newBuilder().setSuccess(false).setErrMsg(errorMessage).build();
    23. } finally {
    24. readLock.unlock();
    25. }
    26. }
    27. private void updateInstanceMetadata(MetadataOperation op) {
    28. Service service = Service.newService(op.getNamespace(), op.getGroup(), op.getServiceName());
    29. service = ServiceManager.getInstance().getSingleton(service);
    30. namingMetadataManager.updateInstanceMetadata(service, op.getTag(), op.getMetadata());
    31. NotifyCenter.publishEvent(new ServiceEvent.ServiceChangedEvent(service, true));
    32. }
    33. private void deleteInstanceMetadata(MetadataOperation op) {
    34. Service service = Service.newService(op.getNamespace(), op.getGroup(), op.getServiceName());
    35. service = ServiceManager.getInstance().getSingleton(service);
    36. namingMetadataManager.removeInstanceMetadata(service, op.getTag());
    37. }

    核心在NamingMetadataManager的updateInstanceMetadata中

    1. private ConcurrentMap> instanceMetadataMap;
    2. /**
    3. * Update instance metadata.
    4. *
    5. * @param service service
    6. * @param metadataId instance metadata id
    7. * @param instanceMetadata new instance metadata
    8. */
    9. public void updateInstanceMetadata(Service service, String metadataId, InstanceMetadata instanceMetadata) {
    10. if (!instanceMetadataMap.containsKey(service)) {
    11. instanceMetadataMap.putIfAbsent(service, new ConcurrentHashMap<>(INITIAL_CAPACITY));
    12. }
    13. instanceMetadataMap.get(service).put(metadataId, instanceMetadata);
    14. }

    其存放客户端实例的对象是一个两层的ConcurrentHashMap,第一层的Map key为service对象,第二层的Map Key为具体的实例Id,value为实例对象。

  • 相关阅读:
    使用爬虫批量下载图片链接并去重
    SMTP发送邮件时抱No appropriate protocol错误分析和解决方案
    centos rpm方式安装jenkins
    如何使用DBeaver连接Hive
    golang 函数式编程库samber/mo使用: Future
    Cache学习(1):常见的程序运行模型&多级Cache存储结构
    如何将AI智能分析与视频监控平台EasyCVR相融合构建监狱安防体系
    centos7篇---安装nvidia-docker
    区域气象-大气化学在线耦合模式(WRF/Chem)在大气环境领域实践技术应用
    java数据结构与算法刷题-----LeetCode104:二叉树的最大深度
  • 原文地址:https://blog.csdn.net/qq_40992849/article/details/133203411