SpringCloud(第 049 篇)Netflix Eureka 源码深入剖析(上)

SpringCloud(第 049 篇)Netflix Eureka 源码深入剖析(上)

一、大致介绍

11、鉴于一些朋友的提问并提议讲解下eureka的源码分析,由此应运而产生的本章节的内容; 22、所以我站在自我的理解角度试着整理了这篇Eureka源码的分析,希望对大家有所帮助; 33、由于篇幅太长不能在一篇里面发布出来,所以拆分了上下篇;

二、基本原理

11、Eureka Server 提供服务注册服务,各个节点启动后,会在Eureka Server中进行注册,这样Eureka Server中的服务注册表中将会存储所有可用服务节点的信息,服务节点的信息可以在界面中直观的看到。 22、Eureka Client 是一个Java 客户端,用于简化与Eureka Server的交互,客户端同时也具备一个内置的、使用轮询负载算法的负载均衡器。 33、在应用启动后,将会向Eureka Server发送心跳(默认周期为30),如果Eureka Server在多个心跳周期没有收到某个节点的心跳,Eureka Server 将会从服务注册表中把这个服务节点移除(默认90)44、Eureka Server之间将会通过复制的方式完成数据的同步; 55、Eureka Client具有缓存的机制,即使所有的Eureka Server 都挂掉的话,客户端依然可以利用缓存中的信息消费其它服务的API

三、EurekaServer 启动流程分析

3.1 跑一下 springms-discovery-eureka 代码,不难发现,我们会看到一些有关 EurekaServer 启动的流程日志;

12017-10-22 18:14:17.635 INFO 5288 --- [ main] o.s.j.e.a.AnnotationMBeanExporter : Located managed bean 'environmentManager': registering with JMX server as MBean [org.springframework.cloud.context.environment:name=environmentManager,type=EnvironmentManager] 22017-10-22 18:14:17.650 INFO 5288 --- [ main] o.s.j.e.a.AnnotationMBeanExporter : Located managed bean 'restartEndpoint': registering with JMX server as MBean [org.springframework.cloud.context.restart:name=restartEndpoint,type=RestartEndpoint] 32017-10-22 18:14:17.661 INFO 5288 --- [ main] o.s.j.e.a.AnnotationMBeanExporter : Located managed bean 'refreshScope': registering with JMX server as MBean [org.springframework.cloud.context.scope.refresh:name=refreshScope,type=RefreshScope] 42017-10-22 18:14:17.674 INFO 5288 --- [ main] o.s.j.e.a.AnnotationMBeanExporter : Located managed bean 'configurationPropertiesRebinder': registering with JMX server as MBean [org.springframework.cloud.context.properties:name=configurationPropertiesRebinder,context=335b5620,type=ConfigurationPropertiesRebinder] 52017-10-22 18:14:17.683 INFO 5288 --- [ main] o.s.j.e.a.AnnotationMBeanExporter : Located managed bean 'refreshEndpoint': registering with JMX server as MBean [org.springframework.cloud.endpoint:name=refreshEndpoint,type=RefreshEndpoint] 62017-10-22 18:14:17.926 INFO 5288 --- [ main] o.s.c.support.DefaultLifecycleProcessor : Starting beans in phase 0 72017-10-22 18:14:17.927 INFO 5288 --- [ main] c.n.e.EurekaDiscoveryClientConfiguration : Registering application unknown with eureka with status UP 82017-10-22 18:14:17.927 INFO 5288 --- [ Thread-10] o.s.c.n.e.server.EurekaServerBootstrap : Setting the eureka configuration.. 92017-10-22 18:14:17.948 INFO 5288 --- [ Thread-10] o.s.c.n.e.server.EurekaServerBootstrap : isAws returned false 102017-10-22 18:14:17.949 INFO 5288 --- [ Thread-10] o.s.c.n.e.server.EurekaServerBootstrap : Initialized server context 112017-10-22 18:14:17.949 INFO 5288 --- [ Thread-10] c.n.e.r.PeerAwareInstanceRegistryImpl : Got 1 instances from neighboring DS node 122017-10-22 18:14:17.949 INFO 5288 --- [ Thread-10] c.n.e.r.PeerAwareInstanceRegistryImpl : Renew threshold is: 1 132017-10-22 18:14:17.949 INFO 5288 --- [ Thread-10] c.n.e.r.PeerAwareInstanceRegistryImpl : Changing status to UP 142017-10-22 18:14:17.958 INFO 5288 --- [ Thread-10] e.s.EurekaServerInitializerConfiguration : Started Eureka Server 152017-10-22 18:14:18.019 INFO 5288 --- [ main] s.b.c.e.t.TomcatEmbeddedServletContainer : Tomcat started on port(s): 8761 (http) 162017-10-22 18:14:18.020 INFO 5288 --- [ main] c.n.e.EurekaDiscoveryClientConfiguration : Updating port to 8761 172017-10-22 18:14:18.023 INFO 5288 --- [ main] c.s.cloud.EurekaServerApplication : Started EurekaServerApplication in 8.299 seconds (JVM running for 8.886) 18【【【【【【 Eureka微服务 】】】】】】已启动. 19 20【分析】:发现有这么一句日志打印“Setting the eureka configuration..”,eureka 开始进行配置,说不定也许就是Eureka Server 流程启动的开 21始呢?我们抱着怀疑的心态进入这行日志打印的EurekaServerBootstrap类去看看。

3.2 进入 EurekaServerBootstrap 类看看,看这个类的名字,见名知意,应该就是 EurekaServer 的启动类了;

1protected void initEurekaEnvironment() throws Exception { 2 log.info("Setting the eureka configuration.."); 3 。。。 45 6【分析一】:我们看到日志在 initEurekaEnvironment 方法中被打印出来,然后我顺着这个方法寻找该方法被调用的地方; 7 8public void contextInitialized(ServletContext context) { 9 try { 10 initEurekaEnvironment(); 11 initEurekaServerContext(); 12 13 context.setAttribute(EurekaServerContext.class.getName(), this.serverContext); 14 } 15 catch (Throwable e) { 16 log.error("Cannot bootstrap eureka server :", e); 17 throw new RuntimeException("Cannot bootstrap eureka server :", e); 18 } 19} 20 21【分析二】:接着发现 contextInitialized 这个方法里面调用了 initEurekaEnvironment 方法,接着我们再往上层寻找被调用的地方; 22 23【分析三】:接着我们看到 EurekaServerInitializerConfiguration 类中有个 start 方法,该方法创建了一个线程来后台执行 EurekaServer 的初始化流程;

3.3 进入 EurekaServerInitializerConfiguration 方法,看看这个所谓的 EurekaServer 初始化配置做了哪些事情?

1@Override 2public void start() { // 打上断点 3 new Thread(new Runnable() { 4 @Override 5 public void run() { 6 try { 7 //TODO: is this class even needed now? 8 eurekaServerBootstrap.contextInitialized(EurekaServerInitializerConfiguration.this.servletContext); 9 log.info("Started Eureka Server"); 10 11 publish(new EurekaRegistryAvailableEvent(getEurekaServerConfig())); 12 EurekaServerInitializerConfiguration.this.running = true; 13 publish(new EurekaServerStartedEvent(getEurekaServerConfig())); 14 } 15 catch (Exception ex) { 16 // Help! 17 log.error("Could not initialize Eureka servlet context", ex); 18 } 19 } 20 }).start(); 21} 22 23【分析一】:看到 log.info("Started Eureka Server"); 这行代码,相信大家已经释然了,这里就是所谓的启动了 EurekaServer 了,其实也就是 24eurekaServerBootstrap.contextInitialized(EurekaServerInitializerConfiguration.this.servletContext) 初始化了一些我们未知的东西; 25 26【分析二】:当打印完启动Eureka Server日志后,调用了两次 publish 方法,该方法最终调用的是 this.applicationContext.publishEvent 27(event) 方法,目的是利用Spring中ApplicationContext对事件传递性质,事件发布者(applicationContext)来发布事件(event),但是缺少的是监听 28者,其实你仔细搜索下代码,发现好像没有地方对 EurekaServerStartedEvent、EurekaRegistryAvailableEvent 进行监听,奇了怪了,这是咋了呢? 29 30【分析三】:然后找到 EurekaServerStartedEvent 所在的目录下,EurekaInstanceCanceledEvent、EurekaInstanceRegisteredEvent、 31EurekaInstanceRenewedEvent、EurekaRegistryAvailableEvent、EurekaServerStartedEvent 有这么几个事件的类,服务下线事件、服务注册事 32件、服务续约事件、注册中心启动事件、Eureka Server启动事件,这么几个事件都没有被监听,那么我们是不是给添加上监听,是不是就可以了呢?像这样 33 @EventListener public void listen(EurekaInstanceCanceledEvent event) { 。。。处下线逻辑 },添加 EventListener 监听注解,就可 34以在我们自己的代码逻辑中收到这个事件的回调了,所以想想SpringCloud还是挺机制的,提供回调接口让我们自己实现自己的业务逻辑,真心不错; 35 36【分析四】:那么反过来想想,为啥会无缘无故 start 方法就被调用了呢?那么反向继续向上找调用 start 方法的地方,结果找到了 37DefaultLifecycleProcessor类的doStart方法调用了 bean.start(); 这么一段代码;

3.4 进入 DefaultLifecycleProcessor 类看看,这个 EurekaServerInitializerConfiguration.start 方法是如何被触发的?

1private void doStart(Map<String, ? extends Lifecycle> lifecycleBeans, String beanName, boolean autoStartupOnly) { 2 // 打上断点 3 Lifecycle bean = lifecycleBeans.remove(beanName); 4 if (bean != null && !this.equals(bean)) { 5 String[] dependenciesForBean = this.beanFactory.getDependenciesForBean(beanName); 6 for (String dependency : dependenciesForBean) { 7 doStart(lifecycleBeans, dependency, autoStartupOnly); 8 } 9 if (!bean.isRunning() && 10 (!autoStartupOnly || !(bean instanceof SmartLifecycle) || ((SmartLifecycle) bean).isAutoStartup())) { 11 if (logger.isDebugEnabled()) { 12 logger.debug("Starting bean '" + beanName + "' of type [" + bean.getClass() + "]"); 13 } 14 try { 15 bean.start(); 16 } 17 catch (Throwable ex) { 18 throw new ApplicationContextException("Failed to start bean '" + beanName + "'", ex); 19 } 20 if (logger.isDebugEnabled()) { 21 logger.debug("Successfully started bean '" + beanName + "'"); 22 } 23 } 24 } 25} 26 27【分析一】:看到在 bean.isRunning 等一系列状态的判断下才去调用 bean.start() 方法的,我们再往上寻找被调用地方; 28 29public void start() { 30 // 打上断点 31 if (this.members.isEmpty()) { 32 return; 33 } 34 if (logger.isInfoEnabled()) { 35 logger.info("Starting beans in phase " + this.phase); 36 } 37 Collections.sort(this.members); 38 for (LifecycleGroupMember member : this.members) { 39 if (this.lifecycleBeans.containsKey(member.name)) { 40 doStart(this.lifecycleBeans, member.name, this.autoStartupOnly); 41 } 42 } 43} 44 45【分析二】:该类是DefaultLifecycleProcessor中内部类LifecycleGroup的一个方法,再往上寻找调用方; 46 47private void startBeans(boolean autoStartupOnly) { 48 Map<String, Lifecycle> lifecycleBeans = getLifecycleBeans(); 49 Map<Integer, LifecycleGroup> phases = new HashMap<Integer, LifecycleGroup>(); 50 for (Map.Entry<String, ? extends Lifecycle> entry : lifecycleBeans.entrySet()) { 51 Lifecycle bean = entry.getValue(); 52 if (!autoStartupOnly || (bean instanceof SmartLifecycle && ((SmartLifecycle) bean).isAutoStartup())) { 53 int phase = getPhase(bean); 54 LifecycleGroup group = phases.get(phase); 55 if (group == null) { 56 group = new LifecycleGroup(phase, this.timeoutPerShutdownPhase, lifecycleBeans, autoStartupOnly); 57 phases.put(phase, group); 58 } 59 group.add(entry.getKey(), bean); 60 } 61 } 62 if (phases.size() > 0) { 63 List<Integer> keys = new ArrayList<Integer>(phases.keySet()); 64 Collections.sort(keys); 65 for (Integer key : keys) { 66 phases.get(key).start(); 67 } 68 } 69} 70 71【分析三】:startBeans 属于 DefaultLifecycleProcessor 类的一个私有方法,startBeans 方法第一行就是获取 getLifecycleBeans() 生命周期 72Bean对象,由此可见似乎 Eureka Server 之所以会被启动,是不是实现了某个接口或者重写了某个方法,才会导致由于容易在初始化的过程中因调用某些特 73殊方法或者某些类才启动的,因此我们回头去看看 EurekaServerInitializerConfiguration 这个类; 74 75【分析四】:结果发现 EurekaServerInitializerConfiguration 这个类实现了 SmartLifecycle 这么样的一个接口,而 SmartLifecycle 接口又继 76承了 Lifecycle 生命周期接口类,所以真想已经重见天日了,原来是实现了 Lifecycle 这样的一个接口,然后实现了 start 方法,因此 Eureka 77Server 就这么稀里糊涂的就被莫名其妙的启动起来了?

3.5 到这里难道就真的完了么?难道Eureka Server启动就干这么点点事情?不可能吧?

1【分析一】:我们之前仅仅只是通过了日志来逆向分析,但是我们是不是忘了我们本应该标志是Eureka Server的这个注解了呢?没错,我们在分析的过程中 2已经将 @EnableEurekaServer 这个注解遗忘了,那么我们现在先回到这个注解类来看看;

3.6 进入 EnableEurekaServer 类,看看究竟干了啥?

1@Target(ElementType.TYPE) 2@Retention(RetentionPolicy.RUNTIME) 3@Documented 4@Import(EurekaServerConfiguration.class) 5public @interface EnableEurekaServer { 6 7} 8 9【分析一】:我们不难发现 EnableEurekaServer 类上有个 @Import 注解,引用了一个 class 文件,由此我们进入观察;

3.7 进入 EurekaServerConfiguration 类看看,看名称的话,理解的意思大概就是 Eureka Server 配置类;

1【分析一】:果不其然,这个类有很多 @Bean、@Configuration 注解过的方法,那是不是我们可以认为刚才 3.1~3.4 的推论是不是就是由于被实例化了这么一个 Bean,然后就慢慢的调用到了 start 方法了呢? 2 3【分析二】:搜索 “Bootstrap” 字样,还真发现了有这么一个方法; 4 5@Bean 6public EurekaServerBootstrap eurekaServerBootstrap(PeerAwareInstanceRegistry registry, 7 EurekaServerContext serverContext) { 8 return new EurekaServerBootstrap(this.applicationInfoManager, 9 this.eurekaClientConfig, this.eurekaServerConfig, registry, 10 serverContext); 11} 12 13【分析三】:既然有这么一个 Bean,那么是不是和刚开始顺着日志逆向分析也是有一定道理的,没有这么一个Bean的存在,那么 DefaultLifecycleProcessor.startBeans 方法中 getLifecycleBeans 的这个也就没那么顺畅被找到了呢?不过我的猜想是这样的,本人没有将源码下载下来,将 eurekaServerBootstrap 方法中的 @Bean 注解注释掉试试,不过推理起来也八九不离十,这个疑问悬念就留给大家尝试尝试吧; 14 15【分析四】:既然找到了一个 @Bean 注解过的方法,那我们再找找其他的一些被注解过的方法,比如一些通用全局用的类似词眼,比如 Context,Bean,Init、Server 之类的; 16 17@Bean 18public EurekaServerContext eurekaServerContext(ServerCodecs serverCodecs, 19 PeerAwareInstanceRegistry registry, PeerEurekaNodes peerEurekaNodes) { 20 return new DefaultEurekaServerContext(this.eurekaServerConfig, serverCodecs, 21 registry, peerEurekaNodes, this.applicationInfoManager); 22} 23 24@Bean 25public PeerEurekaNodes peerEurekaNodes(PeerAwareInstanceRegistry registry, 26 ServerCodecs serverCodecs) { 27 return new PeerEurekaNodes(registry, this.eurekaServerConfig, 28 this.eurekaClientConfig, serverCodecs, this.applicationInfoManager); 29} 30 31@Bean 32public PeerAwareInstanceRegistry peerAwareInstanceRegistry( 33 ServerCodecs serverCodecs) { 34 this.eurekaClient.getApplications(); // force initialization 35 return new InstanceRegistry(this.eurekaServerConfig, this.eurekaClientConfig, 36 serverCodecs, this.eurekaClient, 37 this.instanceRegistryProperties.getExpectedNumberOfRenewsPerMin(), 38 this.instanceRegistryProperties.getDefaultOpenForTrafficCount()); 39} 40 41@Bean 42@ConditionalOnProperty(prefix = "eureka.dashboard", name = "enabled", matchIfMissing = true) 43public EurekaController eurekaController() { 44 return new EurekaController(this.applicationInfoManager); 45} 46 47【分析五】:DefaultEurekaServerContext.initialize 初始化了一些东西,现在还不知道干啥用的,先放这里,打上断点; 48 49【分析六】:PeerEurekaNodes.start 方法,又是一个 start 方法,但是该类没有实现任何类,姑且先放这里,打上断点; 50 51【分析七】:InstanceRegistry.register 方法,而且还有几个呢,可能是客户端注册用的,也先放这里,都打上断点,或者将 这个类的所有方法都断点上,断点打完后发现有注册的,有续约的,有注销的; 52 53【分析八】:打完这些断点后,感觉没有思路了,索性就断点跑一把,看看有什么新的发现点;

3.8 停止服务,Debug 跑一下 springms-discovery-eureka 代码;

1【分析一】:DefaultEurekaServerContext.initialize 方法被调用了,证实了刚才想法,EurekaServerConfiguration 不是白写的,还是有它的作用的; 2 3@PostConstruct 4@Override 5public void initialize() throws Exception { 6 logger.info("Initializing ..."); 7 peerEurekaNodes.start(); 8 registry.init(peerEurekaNodes); 9 logger.info("Initialized"); 10} 11 12【分析二】:进入 initialize 方法中 peerEurekaNodes.start(); 13 14public void start() { 15 taskExecutor = Executors.newSingleThreadScheduledExecutor( 16 new ThreadFactory() { 17 @Override 18 public Thread newThread(Runnable r) { 19 Thread thread = new Thread(r, "Eureka-PeerNodesUpdater"); 20 thread.setDaemon(true); 21 return thread; 22 } 23 } 24 ); 25 try { 26 updatePeerEurekaNodes(resolvePeerUrls()); 27 Runnable peersUpdateTask = new Runnable() { 28 @Override 29 public void run() { 30 try { 31 updatePeerEurekaNodes(resolvePeerUrls()); 32 } catch (Throwable e) { 33 logger.error("Cannot update the replica Nodes", e); 34 } 35 36 } 37 }; 38 // 注释:间隔 600000 毫秒,即 10分钟 间隔执行一次服务集群数据同步; 39 taskExecutor.scheduleWithFixedDelay( 40 peersUpdateTask, 41 serverConfig.getPeerEurekaNodesUpdateIntervalMs(), 42 serverConfig.getPeerEurekaNodesUpdateIntervalMs(), 43 TimeUnit.MILLISECONDS 44 ); 45 } catch (Exception e) { 46 throw new IllegalStateException(e); 47 } 48 for (PeerEurekaNode node : peerEurekaNodes) { 49 logger.info("Replica node URL: " + node.getServiceUrl()); 50 } 51} 52 53【分析三】: start 方法中会看到一个定时调度的任务,updatePeerEurekaNodes(resolvePeerUrls()); 间隔 600000 毫秒,即 10分钟 间隔执行一次服务集群数据同步; 54 55【分析四】: 然后断点放走放下走,进入 initialize 方法中 registry.init(peerEurekaNodes); 56 57@Override 58public void init(PeerEurekaNodes peerEurekaNodes) throws Exception { 59 this.numberOfReplicationsLastMin.start(); 60 this.peerEurekaNodes = peerEurekaNodes; 61 // 注释:初始化 Eureka Server 响应缓存,默认缓存时间为30s 62 initializedResponseCache(); 63 // 注释:定时任务,多久重置一下心跳阈值,900000 毫秒,即 15分钟 的间隔时间,会重置心跳阈值 64 scheduleRenewalThresholdUpdateTask(); 65 // 注释:初始化远端注册 66 initRemoteRegionRegistry(); 67 68 try { 69 Monitors.registerObject(this); 70 } catch (Throwable e) { 71 logger.warn("Cannot register the JMX monitor for the InstanceRegistry :", e); 72 } 73} 74 75【分析五】: 缓存也配置好了,定时任务也配置好了,似乎应该没啥了,那么我们把断点放开,看看下一步会走到哪里?

3.9 EurekaServerInitializerConfiguration.start 也进断点了。

1【分析一】:先是 DefaultLifecycleProcessor.doStart 方法进断点,然后才是 EurekaServerInitializerConfiguration.start 方法进断点; 2 3【分析二】:再一次证明刚刚的逆向分析仅仅只是缺了个从头EnableEurekaServer分析罢了,但是最终方法论分析思路还是对的,由于开始分析过这里,然而我们就跳过,继续放开断点向后继续看看; 4

3.10 InstanceRegistry.openForTraffic 也进断点了。

1【分析一】:这不就是我们刚才在 “步骤3.7之分析七” 打的断点么?看下堆栈信息,正是 “步骤3.2之分析一” 中 initEurekaServerContext 方法中有 2这么一句 this.registry.openForTraffic(this.applicationInfoManager, registryCount); 调用到了,因果轮回,代码千变万化,打上断点还有有好处的,结果还是回到了开始日志逆向分析的地方。 3 4【分析二】:进入 super.openForTraffic 方法; 5 6@Override 7public void openForTraffic(ApplicationInfoManager applicationInfoManager, int count) { 8 // Renewals happen every 30 seconds and for a minute it should be a factor of 2. 9 // 注释:每30秒续约一次,那么每分钟续约就是2次,所以才是 count * 2 的结果; 10 this.expectedNumberOfRenewsPerMin = count * 2; 11 this.numberOfRenewsPerMinThreshold = 12 (int) (this.expectedNumberOfRenewsPerMin * serverConfig.getRenewalPercentThreshold()); 13 logger.info("Got " + count + " instances from neighboring DS node"); 14 logger.info("Renew threshold is: " + numberOfRenewsPerMinThreshold); 15 this.startupTime = System.currentTimeMillis(); 16 if (count > 0) { 17 this.peerInstancesTransferEmptyOnStartup = false; 18 } 19 DataCenterInfo.Name selfName = applicationInfoManager.getInfo().getDataCenterInfo().getName(); 20 boolean isAws = Name.Amazon == selfName; 21 if (isAws && serverConfig.shouldPrimeAwsReplicaConnections()) { 22 logger.info("Priming AWS connections for all replicas.."); 23 primeAwsReplicas(applicationInfoManager); 24 } 25 logger.info("Changing status to UP"); 26 // 注释:修改 Eureka Server 为上电状态,就是说设置 Eureka Server 已经处于活跃状态了,那就是意味着 EurekaServer 基本上说可以正常使用了; 27 applicationInfoManager.setInstanceStatus(InstanceStatus.UP); 28 // 注释:定时任务,60000 毫秒,即 1分钟 的间隔时间,Eureke Server定期进行失效节点的清理 29 super.postInit(); 30} 31 32【分析三】:这里主要设置了服务状态,以及开启了定时清理失效节点的定时任务,每分钟扫描一次;

3.11 继续放开断点,来到了日志打印 “main] c.n.e.EurekaDiscoveryClientConfiguration : Updating port to 8761” 的EurekaDiscoveryClientConfiguration 类中 onApplicationEvent 方法。

1@EventListener(EmbeddedServletContainerInitializedEvent.class) 2public void onApplicationEvent(EmbeddedServletContainerInitializedEvent event) { 3 // TODO: take SSL into account when Spring Boot 1.2 is available 4 int localPort = event.getEmbeddedServletContainer().getPort(); 5 if (this.port.get() == 0) { 6 log.info("Updating port to " + localPort); 7 this.port.compareAndSet(0, localPort); 8 start(); 9 } 10} 11 12【分析一】:设置端口,当看到 Updating port to 8761 这样的日志打印出来的话,说明 Eureka Server 整个启动也就差不多Over了。现在回头看看, 13发现分析了不少的方法和流程,有种感觉被掏空的感觉了。

3.12 总结 EurekaServer 启动时候大概干了哪些事情?

11、初始化Eureka环境,Eureka上下文; 22、初始化EurekaServer的缓存 33、启动了一些定时任务,比如充值心跳阈值定时任务,清理失效节点定时任务; 44、更新EurekaServer上电状态,更新EurekaServer端口; 5 6虽然我从列举的流程里面大概总结了这么几点,但是还是有些是我没关注到的,如果大家有关注到的,可以和我共同讨论分析分析。

四、EurekaServer 处理服务注册、集群数据复制

4.1 EurekaClient 是如何注册到 EurekaServer 的?

1【分析一】:由于我们刚才在 org.springframework.cloud.netflix.eureka.server.InstanceRegistry 的每个方法都打了一个断点,而且现在 2EurekaServer 已经处于 Debug 运行状态,那么我们就随便找一个被 @EnableEurekaClient 的微服务启动试试,要么就拿 springms-provider-user 3微服务来试试吧,直接 Run。 4 5【分析二】:猜测,如果如我们分析所想,当 springms-provider-user 启动后,就一定会调用注册register方法,那么就接着往下看,拭目以待;

4.2 InstanceRegistry.register(final InstanceInfo info, final boolean isReplication) 方法进断点了。

【分析一】:由于 InstanceRegistry.register 是我们刚刚打断点的地方,那么我们顺着堆栈信息往上看,原来是 ApplicationResource.addInstance 方法被调用了,那么我们就看看 addInstance 这个方法,并在 addInstance 这里打上断点;接着我们重新杀死 springms-provider-user 服务,然后再重启 springms-provider-user 服务;

4.2 断点再次来到了 ApplicationResource 类,这个类呢,主要是处理接收 Http 的服务请求。

1@POST 2@Consumes({"application/json", "application/xml"}) 3public Response addInstance(InstanceInfo info, 4 @HeaderParam(PeerEurekaNode.HEADER_REPLICATION) String isReplication) { 5 logger.debug("Registering instance {} (replication={})", info.getId(), isReplication); 6 // validate that the instanceinfo contains all the necessary required fields 7 if (isBlank(info.getId())) { 8 return Response.status(400).entity("Missing instanceId").build(); 9 } else if (isBlank(info.getHostName())) { 10 return Response.status(400).entity("Missing hostname").build(); 11 } else if (isBlank(info.getAppName())) { 12 return Response.status(400).entity("Missing appName").build(); 13 } else if (!appName.equals(info.getAppName())) { 14 return Response.status(400).entity("Mismatched appName, expecting " + appName + " but was " + info.getAppName()).build(); 15 } else if (info.getDataCenterInfo() == null) { 16 return Response.status(400).entity("Missing dataCenterInfo").build(); 17 } else if (info.getDataCenterInfo().getName() == null) { 18 return Response.status(400).entity("Missing dataCenterInfo Name").build(); 19 } 20 21 // handle cases where clients may be registering with bad DataCenterInfo with missing data 22 DataCenterInfo dataCenterInfo = info.getDataCenterInfo(); 23 if (dataCenterInfo instanceof UniqueIdentifier) { 24 String dataCenterInfoId = ((UniqueIdentifier) dataCenterInfo).getId(); 25 if (isBlank(dataCenterInfoId)) { 26 boolean experimental = "true".equalsIgnoreCase(serverConfig.getExperimental("registration.validation.dataCenterInfoId")); 27 if (experimental) { 28 String entity = "DataCenterInfo of type " + dataCenterInfo.getClass() + " must contain a valid id"; 29 return Response.status(400).entity(entity).build(); 30 } else if (dataCenterInfo instanceof AmazonInfo) { 31 AmazonInfo amazonInfo = (AmazonInfo) dataCenterInfo; 32 String effectiveId = amazonInfo.get(AmazonInfo.MetaDataKey.instanceId); 33 if (effectiveId == null) { 34 amazonInfo.getMetadata().put(AmazonInfo.MetaDataKey.instanceId.getName(), info.getId()); 35 } 36 } else { 37 logger.warn("Registering DataCenterInfo of type {} without an appropriate id", dataCenterInfo.getClass()); 38 } 39 } 40 } 41 42 registry.register(info, "true".equals(isReplication)); 43 return Response.status(204).build(); // 204 to be backwards compatible 44} 45 46【分析一】:这里的写法貌似看起来和我们之前 ControllerRESTFUL 写法有点不一样,仔细一看,原来是Jersey RESTful 框架,是一个产品级的 47RESTful service 和 client 框架。与Struts类似,它同样可以和hibernate,spring框架整合。 48 49【分析二】:紧接着,我们看到 registry.register(info, "true".equals(isReplication)); 这么一段代码,注册啊,原来EurekaClient客户端启 50动后会调用会通过Http(s)请求,直接调到 ApplicationResource.addInstance 方法,那么总算明白了,只要是和注册有关的,都会调用这个方法。 51 52【分析三】:接着我们深入 registry.register(info, "true".equals(isReplication)) 查看; 53 54@Override 55public void register(final InstanceInfo info, final boolean isReplication) { 56 handleRegistration(info, resolveInstanceLeaseDuration(info), isReplication); 57 super.register(info, isReplication); 58} 59 60【分析四】:看看 handleRegistration(info, resolveInstanceLeaseDuration(info), isReplication) 方法; 61 62private void handleRegistration(InstanceInfo info, int leaseDuration, 63 boolean isReplication) { 64 log("register " + info.getAppName() + ", vip " + info.getVIPAddress() 65 + ", leaseDuration " + leaseDuration + ", isReplication " 66 + isReplication); 67 publishEvent(new EurekaInstanceRegisteredEvent(this, info, leaseDuration, 68 isReplication)); 69} 70 71【分析五】:该方法仅仅只是打了一个日志,然后通过 ApplicationContext 发布了一个事件 EurekaInstanceRegisteredEvent 服务注册事件,正如 72“步骤3.3之分析三” 所提到的,用户可以给 EurekaInstanceRegisteredEvent 添加监听事件,那么用户就可以在此刻实现自己想要的一些业务逻辑。 73然后我们再来看看 super.register(info, isReplication) 方法,该方法是 InstanceRegistry 的父类 PeerAwareInstanceRegistryImpl 的方法。

4.3 进入 PeerAwareInstanceRegistryImpl 类的 register(final InstanceInfo info, final boolean isReplication) 方法;

1@Override 2public void register(final InstanceInfo info, final boolean isReplication) { 3 // 注释:续约时间,默认时间是常量值 90 秒 4 int leaseDuration = Lease.DEFAULT_DURATION_IN_SECS; 5 // 注释:续约时间,当然也可以从配置文件中取出来,所以说续约时间值也是可以让我们自己自定义配置的 6 if (info.getLeaseInfo() != null && info.getLeaseInfo().getDurationInSecs() > 0) { 7 leaseDuration = info.getLeaseInfo().getDurationInSecs(); 8 } 9 // 注释:将注册方的信息写入 EurekaServer 的注册表,父类为 AbstractInstanceRegistry 10 super.register(info, leaseDuration, isReplication); 11 // 注释:EurekaServer 节点之间的数据同步,复制到其他Peer 12 replicateToPeers(Action.Register, info.getAppName(), info.getId(), info, null, isReplication); 13} 14 15【分析一】:进入 super.register(info, leaseDuration, isReplication) 看看是如何写入 EurekaServer 的注册表的,即进入 AbstractInstanceRegistry.register(InstanceInfo registrant, int leaseDuration, boolean isReplication) 方法。 16 17public void register(InstanceInfo registrant, int leaseDuration, boolean isReplication) { 18 try { 19 read.lock(); 20 // 注释:registry 这个变量,就是我们所谓的注册表,注册表是保存在内存中的; 21 Map<String, Lease<InstanceInfo>> gMap = registry.get(registrant.getAppName()); 22 REGISTER.increment(isReplication); 23 if (gMap == null) { 24 final ConcurrentHashMap<String, Lease<InstanceInfo>> gNewMap = new ConcurrentHashMap<String, Lease<InstanceInfo>>(); 25 gMap = registry.putIfAbsent(registrant.getAppName(), gNewMap); 26 if (gMap == null) { 27 gMap = gNewMap; 28 } 29 } 30 Lease<InstanceInfo> existingLease = gMap.get(registrant.getId()); 31 // Retain the last dirty timestamp without overwriting it, if there is already a lease 32 if (existingLease != null && (existingLease.getHolder() != null)) { 33 Long existingLastDirtyTimestamp = existingLease.getHolder().getLastDirtyTimestamp(); 34 Long registrationLastDirtyTimestamp = registrant.getLastDirtyTimestamp(); 35 logger.debug("Existing lease found (existing={}, provided={}", existingLastDirtyTimestamp, registrationLastDirtyTimestamp); 36 if (existingLastDirtyTimestamp > registrationLastDirtyTimestamp) { 37 logger.warn("There is an existing lease and the existing lease's dirty timestamp {} is greater" + 38 " than the one that is being registered {}", existingLastDirtyTimestamp, registrationLastDirtyTimestamp); 39 logger.warn("Using the existing instanceInfo instead of the new instanceInfo as the registrant"); 40 registrant = existingLease.getHolder(); 41 } 42 } else { 43 // The lease does not exist and hence it is a new registration 44 synchronized (lock) { 45 if (this.expectedNumberOfRenewsPerMin > 0) { 46 // Since the client wants to cancel it, reduce the threshold 47 // (1 48 // for 30 seconds, 2 for a minute) 49 this.expectedNumberOfRenewsPerMin = this.expectedNumberOfRenewsPerMin + 2; 50 this.numberOfRenewsPerMinThreshold = 51 (int) (this.expectedNumberOfRenewsPerMin * serverConfig.getRenewalPercentThreshold()); 52 } 53 } 54 logger.debug("No previous lease information found; it is new registration"); 55 } 56 Lease<InstanceInfo> lease = new Lease<InstanceInfo>(registrant, leaseDuration); 57 if (existingLease != null) { 58 lease.setServiceUpTimestamp(existingLease.getServiceUpTimestamp()); 59 } 60 gMap.put(registrant.getId(), lease); 61 synchronized (recentRegisteredQueue) { 62 recentRegisteredQueue.add(new Pair<Long, String>( 63 System.currentTimeMillis(), 64 registrant.getAppName() + "(" + registrant.getId() + ")")); 65 } 66 // This is where the initial state transfer of overridden status happens 67 if (!InstanceStatus.UNKNOWN.equals(registrant.getOverriddenStatus())) { 68 logger.debug("Found overridden status {} for instance {}. Checking to see if needs to be add to the " 69 + "overrides", registrant.getOverriddenStatus(), registrant.getId()); 70 if (!overriddenInstanceStatusMap.containsKey(registrant.getId())) { 71 logger.info("Not found overridden id {} and hence adding it", registrant.getId()); 72 overriddenInstanceStatusMap.put(registrant.getId(), registrant.getOverriddenStatus()); 73 } 74 } 75 InstanceStatus overriddenStatusFromMap = overriddenInstanceStatusMap.get(registrant.getId()); 76 if (overriddenStatusFromMap != null) { 77 logger.info("Storing overridden status {} from map", overriddenStatusFromMap); 78 registrant.setOverriddenStatus(overriddenStatusFromMap); 79 } 80 81 // Set the status based on the overridden status rules 82 InstanceStatus overriddenInstanceStatus = getOverriddenInstanceStatus(registrant, existingLease, isReplication); 83 registrant.setStatusWithoutDirty(overriddenInstanceStatus); 84 85 // If the lease is registered with UP status, set lease service up timestamp 86 if (InstanceStatus.UP.equals(registrant.getStatus())) { 87 lease.serviceUp(); 88 } 89 registrant.setActionType(ActionType.ADDED); 90 recentlyChangedQueue.add(new RecentlyChangedItem(lease)); 91 registrant.setLastUpdatedTimestamp(); 92 invalidateCache(registrant.getAppName(), registrant.getVIPAddress(), registrant.getSecureVipAddress()); 93 logger.info("Registered instance {}/{} with status {} (replication={})", 94 registrant.getAppName(), registrant.getId(), registrant.getStatus(), isReplication); 95 } finally { 96 read.unlock(); 97 } 98} 99 100【分析二】:发现这个方法有点长,大致阅读,主要更新了注册表的时间之外,还更新了缓存等其它东西,大家有兴趣的可以深究阅读该方法;

4.4 跳出来我们接着看上面的 replicateToPeers(Action.Register, info.getAppName(), info.getId(), info, null, isReplication) 的这个方法。

1private void replicateToPeers(Action action, String appName, String id, 2 InstanceInfo info /* optional */, 3 InstanceStatus newStatus /* optional */, boolean isReplication) { 4 Stopwatch tracer = action.getTimer().start(); 5 try { 6 if (isReplication) { 7 numberOfReplicationsLastMin.increment(); 8 } 9 // If it is a replication already, do not replicate again as this will create a poison replication 10 // 注释:如果已经复制过,就不再复制 11 if (peerEurekaNodes == Collections.EMPTY_LIST || isReplication) { 12 return; 13 } 14 15 // 遍历Eureka Server集群中的所有节点,进行复制操作 16 for (final PeerEurekaNode node : peerEurekaNodes.getPeerEurekaNodes()) { 17 // If the url represents this host, do not replicate to yourself. 18 if (peerEurekaNodes.isThisMyUrl(node.getServiceUrl())) { 19 continue; 20 } 21 // 没有复制过,遍历Eureka Server集群中的node节点,依次操作,包括取消、注册、心跳、状态更新等。 22 replicateInstanceActionsToPeers(action, appName, id, info, newStatus, node); 23 } 24 } finally { 25 tracer.stop(); 26 } 27} 28 29【分析一】:走到这里,我不难理解,每当有注册请求,首先更新 EurekaServer 的注册表,然后再将信息同步到其它EurekaServer的节点上去; 30 31【分析二】:接下来我们看看 node 节点是如何进行复制操作的,进入 replicateInstanceActionsToPeers 方法。 32 33private void replicateInstanceActionsToPeers(Action action, String appName, 34 String id, InstanceInfo info, InstanceStatus newStatus, 35 PeerEurekaNode node) { 36 try { 37 InstanceInfo infoFromRegistry = null; 38 CurrentRequestVersion.set(Version.V2); 39 switch (action) { 40 case Cancel: 41 node.cancel(appName, id); 42 break; 43 case Heartbeat: 44 InstanceStatus overriddenStatus = overriddenInstanceStatusMap.get(id); 45 infoFromRegistry = getInstanceByAppAndId(appName, id, false); 46 node.heartbeat(appName, id, infoFromRegistry, overriddenStatus, false); 47 break; 48 case Register: 49 node.register(info); 50 break; 51 case StatusUpdate: 52 infoFromRegistry = getInstanceByAppAndId(appName, id, false); 53 node.statusUpdate(appName, id, newStatus, infoFromRegistry); 54 break; 55 case DeleteStatusOverride: 56 infoFromRegistry = getInstanceByAppAndId(appName, id, false); 57 node.deleteStatusOverride(appName, id, infoFromRegistry); 58 break; 59 } 60 } catch (Throwable t) { 61 logger.error("Cannot replicate information to {} for action {}", node.getServiceUrl(), action.name(), t); 62 } 63} 64 65【分析三】:节点之间的复制状态操作,都在这里体现的淋漓尽致,那么我们就拿 Register 类型 node.register(info) 来看,我们来看看 node 究竟是 66如何做到同步信息的,进入 node.register(info) 方法看看;

4.5 进入 PeerEurekaNode.register(final InstanceInfo info) 方法,一窥究竟如何同步数据。

1public void register(final InstanceInfo info) throws Exception { 2 // 注释:任务过期时间给任务分发器处理,默认时间偏移当前时间 30秒 3 long expiryTime = System.currentTimeMillis() + getLeaseRenewalOf(info); 4 batchingDispatcher.process( 5 taskId("register", info), 6 new InstanceReplicationTask(targetHost, Action.Register, info, null, true) { 7 public EurekaHttpResponse<Void> execute() { 8 return replicationClient.register(info); 9 } 10 }, 11 expiryTime 12 ); 13} 14 15【分析一】:这里涉及到了 Eureka 的任务批处理,通常情况下Peer之间的同步需要调用多次,如果EurekaServer一多的话,那么将会有很多http请求,所 16以自然而然的孕育出了任务批处理,但是也在一定程度上导致了注册和下线的一些延迟,突出优势的同时也势必会造成一些劣势,但是这些延迟情况还是能符合 17常理在容忍范围之内的。 18 19【分析二】:在 expiryTime 超时时间之内,批次处理要做的事情就是合并任务为一个List,然后发送请求的时候,将这个批次List直接打包发送请求出去,这样的话,在这个批次的List里面,可能包含取消、注册、心跳、状态等一系列状态的集合List。 20 21【分析三】:我们再接着看源码,batchingDispatcher.process 这么一调用,然后我们就直接看这个 TaskDispatchers.createBatchingTaskDispatcher 方法。 22 23public static <ID, T> TaskDispatcher<ID, T> createBatchingTaskDispatcher(String id, 24 int maxBufferSize, 25 int workloadSize, 26 int workerCount, 27 long maxBatchingDelay, 28 long congestionRetryDelayMs, 29 long networkFailureRetryMs, 30 TaskProcessor<T> taskProcessor) { 31 final AcceptorExecutor<ID, T> acceptorExecutor = new AcceptorExecutor<>( 32 id, maxBufferSize, workloadSize, maxBatchingDelay, congestionRetryDelayMs, networkFailureRetryMs 33 ); 34 final TaskExecutors<ID, T> taskExecutor = TaskExecutors.batchExecutors(id, workerCount, taskProcessor, acceptorExecutor); 35 return new TaskDispatcher<ID, T>() { 36 @Override 37 public void process(ID id, T task, long expiryTime) { 38 acceptorExecutor.process(id, task, expiryTime); 39 } 40 41 @Override 42 public void shutdown() { 43 acceptorExecutor.shutdown(); 44 taskExecutor.shutdown(); 45 } 46 }; 47 } 48 49【分析四】:这里的 process 方法会将任务添加到队列中,有入队列自然有出队列,具体怎么取任务,我就不一一给大家讲解了,我就讲讲最后是怎么触发任务的。进入 final TaskExecutors<ID, T> taskExecutor = TaskExecutors.batchExecutors(id, workerCount, taskProcessor, acceptorExecutor) 这句代码的 TaskExecutors.batchExecutors 方法。 50 51static <ID, T> TaskExecutors<ID, T> batchExecutors(final String name, 52 int workerCount, 53 final TaskProcessor<T> processor, 54 final AcceptorExecutor<ID, T> acceptorExecutor) { 55 final AtomicBoolean isShutdown = new AtomicBoolean(); 56 final TaskExecutorMetrics metrics = new TaskExecutorMetrics(name); 57 return new TaskExecutors<>(new WorkerRunnableFactory<ID, T>() { 58 @Override 59 public WorkerRunnable<ID, T> create(int idx) { 60 return new BatchWorkerRunnable<>("TaskBatchingWorker-" +name + '-' + idx, isShutdown, metrics, processor, acceptorExecutor); 61 } 62 }, workerCount, isShutdown); 63} 64 65【分析五】:我们发现 TaskExecutors 类中的 batchExecutors 这个静态方法,有个 BatchWorkerRunnable 返回的实现类,因此我们再次进入 BatchWorkerRunnable 类看看究竟,而且既然是 Runnable,那么势必会有 run 方法。 66 67@Override 68public void run() { 69 try { 70 while (!isShutdown.get()) { 71 // 注释:获取信号量释放 batchWorkRequests.release(),返回任务集合列表 72 List<TaskHolder<ID, T>> holders = getWork(); 73 metrics.registerExpiryTimes(holders); 74 75 List<T> tasks = getTasksOf(holders); 76 // 注释:将批量任务打包请求Peer节点 77 ProcessingResult result = processor.process(tasks); 78 switch (result) { 79 case Success: 80 break; 81 case Congestion: 82 case TransientError: 83 taskDispatcher.reprocess(holders, result); 84 break; 85 case PermanentError: 86 logger.warn("Discarding {} tasks of {} due to permanent error", holders.size(), workerName); 87 } 88 metrics.registerTaskResult(result, tasks.size()); 89 } 90 } catch (InterruptedException e) { 91 // Ignore 92 } catch (Throwable e) { 93 // Safe-guard, so we never exit this loop in an uncontrolled way. 94 logger.warn("Discovery WorkerThread error", e); 95 } 96} 97 98【分析六】:这就是我们 BatchWorkerRunnable 类的 run 方法,这里面首先要获取信号量释放,才能获得任务集合,一旦获取到了任务集合的话,那么就直接调用 processor.process(tasks) 方法请求 Peer 节点同步数据,接下来我们看看 ReplicationTaskProcessor.process 方法; 99 100@Override 101public ProcessingResult process(List<ReplicationTask> tasks) { 102 ReplicationList list = createReplicationListOf(tasks); 103 try { 104 // 注释:这里通过 JerseyReplicationClient 客户端对象直接发送list请求数据 105 EurekaHttpResponse<ReplicationListResponse> response = replicationClient.submitBatchUpdates(list); 106 int statusCode = response.getStatusCode(); 107 if (!isSuccess(statusCode)) { 108 if (statusCode == 503) { 109 logger.warn("Server busy (503) HTTP status code received from the peer {}; rescheduling tasks after delay", peerId); 110 return ProcessingResult.Congestion; 111 } else { 112 // Unexpected error returned from the server. This should ideally never happen. 113 logger.error("Batch update failure with HTTP status code {}; discarding {} replication tasks", statusCode, tasks.size()); 114 return ProcessingResult.PermanentError; 115 } 116 } else { 117 handleBatchResponse(tasks, response.getEntity().getResponseList()); 118 } 119 } catch (Throwable e) { 120 if (isNetworkConnectException(e)) { 121 logNetworkErrorSample(null, e); 122 return ProcessingResult.TransientError; 123 } else { 124 logger.error("Not re-trying this exception because it does not seem to be a network exception", e); 125 return ProcessingResult.PermanentError; 126 } 127 } 128 return ProcessingResult.Success; 129} 130 131【分析七】:感觉快要见到真相了,所以我们迫不及待的进入 JerseyReplicationClient.submitBatchUpdates(ReplicationList replicationList) 方法一窥究竟。 132 133@Override 134public EurekaHttpResponse<ReplicationListResponse> submitBatchUpdates(ReplicationList replicationList) { 135 ClientResponse response = null; 136 try { 137 response = jerseyApacheClient.resource(serviceUrl) 138 // 注释:这才是重点,请求目的相对路径,peerreplication/batch/ 139 .path(PeerEurekaNode.BATCH_URL_PATH) 140 .accept(MediaType.APPLICATION_JSON_TYPE) 141 .type(MediaType.APPLICATION_JSON_TYPE) 142 .post(ClientResponse.class, replicationList); 143 if (!isSuccess(response.getStatus())) { 144 return anEurekaHttpResponse(response.getStatus(), ReplicationListResponse.class).build(); 145 } 146 ReplicationListResponse batchResponse = response.getEntity(ReplicationListResponse.class); 147 return anEurekaHttpResponse(response.getStatus(), batchResponse).type(MediaType.APPLICATION_JSON_TYPE).build(); 148 } finally { 149 if (response != null) { 150 response.close(); 151 } 152 } 153} 154 155【分析八】:看到了相对路径地址,我们搜索下"batch"这样的字符串看看有没有对应的接收方法或者被@Path注解进入的;在 eureka-core-1.4.12.jar 这个包下面,果然搜到到了 @Path("batch") 这样的字样,直接进入,发现这是 PeerReplicationResource 类的方法 batchReplication,我们进入这方法看看。 156 157@Path("batch") 158@POST 159public Response batchReplication(ReplicationList replicationList) { 160 try { 161 ReplicationListResponse batchResponse = new ReplicationListResponse(); 162 // 注释:这里将收到的任务列表,依次循环解析处理,主要核心方法在 dispatch 方法中。 163 for (ReplicationInstance instanceInfo : replicationList.getReplicationList()) { 164 try { 165 batchResponse.addResponse(dispatch(instanceInfo)); 166 } catch (Exception e) { 167 batchResponse.addResponse(new ReplicationInstanceResponse(Status.INTERNAL_SERVER_ERROR.getStatusCode(), null)); 168 logger.error(instanceInfo.getAction() + " request processing failed for batch item " 169 + instanceInfo.getAppName() + '/' + instanceInfo.getId(), e); 170 } 171 } 172 return Response.ok(batchResponse).build(); 173 } catch (Throwable e) { 174 logger.error("Cannot execute batch Request", e); 175 return Response.status(Status.INTERNAL_SERVER_ERROR).build(); 176 } 177} 178 179【分析九】:看到了循环一次遍历任务进行处理,不知不觉觉得心花怒放,胜利的重点马上就要到来了,我们进入 PeerReplicationResource.dispatch 方法看看。 180 181private ReplicationInstanceResponse dispatch(ReplicationInstance instanceInfo) { 182 ApplicationResource applicationResource = createApplicationResource(instanceInfo); 183 InstanceResource resource = createInstanceResource(instanceInfo, applicationResource); 184 185 String lastDirtyTimestamp = toString(instanceInfo.getLastDirtyTimestamp()); 186 String overriddenStatus = toString(instanceInfo.getOverriddenStatus()); 187 String instanceStatus = toString(instanceInfo.getStatus()); 188 189 Builder singleResponseBuilder = new Builder(); 190 switch (instanceInfo.getAction()) { 191 case Register: 192 singleResponseBuilder = handleRegister(instanceInfo, applicationResource); 193 break; 194 case Heartbeat: 195 singleResponseBuilder = handleHeartbeat(resource, lastDirtyTimestamp, overriddenStatus, instanceStatus); 196 break; 197 case Cancel: 198 singleResponseBuilder = handleCancel(resource); 199 break; 200 case StatusUpdate: 201 singleResponseBuilder = handleStatusUpdate(instanceInfo, resource); 202 break; 203 case DeleteStatusOverride: 204 singleResponseBuilder = handleDeleteStatusOverride(instanceInfo, resource); 205 break; 206 } 207 return singleResponseBuilder.build(); 208} 209 210【分析十】:随便抓一个类型,那我们也拿 Register 类型来看,进入 PeerReplicationResource.handleRegister 看看。 211 212private static Builder handleRegister(ReplicationInstance instanceInfo, ApplicationResource applicationResource) { 213 // 注释:private static final String REPLICATION = "true"; 定义的一个常量值,而且还是回调 ApplicationResource.addInstance 方法 214 applicationResource.addInstance(instanceInfo.getInstanceInfo(), REPLICATION); 215 return new Builder().setStatusCode(Status.OK.getStatusCode()); 216} 217 218【分析十一】:Peer节点的同步旅程终于结束了,最终又回调到了 ApplicationResource.addInstance 这个方法,这个方法在最终是EurekaClient启动后注册调用的方法,然而Peer节点的信息同步也调用了这个方法,仅仅只是通过一个变量 isReplication 为true还是false来判断是否是节点复制。剩下的ApplicationResource.addInstance流程前面已经提到过了,相信大家已经明白了注册的流程是如何扭转的,包括批量任务是如何处理EurekaServer节点之间的信息同步的了。 219

五、EurekaClient 启动流程分析

详见 SpringCloud(第 050 篇)Netflix Eureka 源码深入剖析(下)

六、下载地址

https://gitee.com/ylimhhmily/SpringCloudTutorial.git

SpringCloudTutorial交流QQ群: 235322432

SpringCloudTutorial交流微信群: 微信沟通群二维码图片链接

欢迎关注,您的肯定是对我最大的支持!!!

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )