Spring Cloud Eureka源代码解析(2) EurekaServer 重要缓存解析

我们从EurekaServer的缓存说起,因为缓存是EurekaServer的一切存储形式,并且我们通过对缓存的分析可以搞清楚一些对于EurekaServer的误解。

服务实例向EurekaServer注册,注册信息是放在缓存中。从EurekaServer中获取服务实例列表的时候,也是从缓存获取;但是这个缓存结构比较复杂,并且还有很多定时刷新和定时失效的机制,我们需要仔细分析

首先,从核心的服务注册信息存储的地方说起,并简单介绍下其中的注册,取消机制。

其核心逻辑位于AbstractInstanceRegistry这个类,这个类有如下几个重要的存储:

1private final ConcurrentHashMap<String, Map<String, Lease<InstanceInfo>>> registry 2 = new ConcurrentHashMap<String, Map<String, Lease<InstanceInfo>>>(); 3private ConcurrentLinkedQueue<RecentlyChangedItem> recentlyChangedQueue = new ConcurrentLinkedQueue<RecentlyChangedItem>();

并由如下几个重要的锁控制:

1private final ReentrantReadWriteLock readWriteLock = new ReentrantReadWriteLock(); 2private final Lock read = readWriteLock.readLock(); 3private final Lock write = readWriteLock.writeLock();

其他的field和一些监控还有自我保护机制相关,我们先不管。

有服务实例注册时,会调用register方法,简化后的代码为:

1public void register(InstanceInfo registrant, int leaseDuration, boolean isReplication) { 2 try { 3 //register虽然看上去好像是修改,但是这里用的是读锁,后面会解释 4 read.lock(); 5 //从registry中查看这个app是否存在 6 Map<String, Lease<InstanceInfo>> gMap = registry.get(registrant.getAppName()); 7 //不存在就创建 8 if (gMap == null) { 9 final ConcurrentHashMap<String, Lease<InstanceInfo>> gNewMap = new ConcurrentHashMap<String, Lease<InstanceInfo>>(); 10 gMap = registry.putIfAbsent(registrant.getAppName(), gNewMap); 11 if (gMap == null) { 12 gMap = gNewMap; 13 } 14 } 15 //查看这个app的这个实例是否已存在 16 Lease<InstanceInfo> existingLease = gMap.get(registrant.getId()); 17 18 if (existingLease != null && (existingLease.getHolder() != null)) { 19 //如果已存在,对比时间戳,保留比较新的实例信息...... 20 } else { 21 // 如果不存在,证明是一个新的实例 22 //更新自我保护监控变量的值的代码..... 23 24 } 25 Lease<InstanceInfo> lease = new Lease<InstanceInfo>(registrant, leaseDuration); 26 if (existingLease != null) { 27 lease.setServiceUpTimestamp(existingLease.getServiceUpTimestamp()); 28 } 29 //放入registry 30 gMap.put(registrant.getId(), lease); 31 32 //加入最近修改的记录队列 33 recentlyChangedQueue.add(new RecentlyChangedItem(lease)); 34 //初始化状态,记录时间等相关代码...... 35 36 //主动让Response缓存失效 37 invalidateCache(registrant.getAppName(), registrant.getVIPAddress(), registrant.getSecureVipAddress()); 38 } finally { 39 read.unlock(); 40 } 41}

总结起来,就是主要三件事:
1.将实例注册信息放入或者更新registry
2.将实例注册信息加入最近修改的记录队列
3.主动让Response缓存失效

这个Response缓存,我们稍后就会介绍

当服务实例取消注册时,会调用cancel方法,cacel直接调用internalCancel,这个internalCancel被抽离出来是因为Eureka主动检查evict机制也会调用这个方法。简化后的internalCancel代码为:

1protected boolean internalCancel(String appName, String id, boolean isReplication) { 2 try { 3 //cancel虽然看上去好像是修改,但是这里用的是读锁,后面会解释 4 read.lock(); 5 6 //从registry中剔除这个实例 7 Map<String, Lease<InstanceInfo>> gMap = registry.get(appName); 8 Lease<InstanceInfo> leaseToCancel = null; 9 if (gMap != null) { 10 leaseToCancel = gMap.remove(id); 11 } 12 if (leaseToCancel == null) { 13 logger.warn("DS: Registry: cancel failed because Lease is not registered for: {}/{}", appName, id); 14 return false; 15 } else { 16 //改变状态,记录状态修改时间等相关代码...... 17 if (instanceInfo != null) { 18 instanceInfo.setActionType(ActionType.DELETED); 19 //加入最近修改的记录队列 20 recentlyChangedQueue.add(new RecentlyChangedItem(leaseToCancel)); 21 } 22 //主动让Response缓存失效 23 invalidateCache(appName, vip, svip); 24 logger.info("Cancelled instance {}/{} (replication={})", appName, id, isReplication); 25 return true; 26 } 27 } finally { 28 read.unlock(); 29 } 30}

总结起来,也是主要三件事:
1.从registry中剔除这个实例
2.将实例注册信息加入最近修改的记录队列
3.主动让Response缓存失效

同时,这个类还会启动两个定时任务,一个是主动失效evict过期应用实例的服务,这里我们先不讨论;另一个就是定时清理最近修改的记录队列的任务:

1Iterator<RecentlyChangedItem> it = recentlyChangedQueue.iterator(); 2 while (it.hasNext()) { 3 if (it.next().getLastUpdateTime() < 4 System.currentTimeMillis() - serverConfig.getRetentionTimeInMSInDeltaQueue()) { 5 it.remove(); 6 } else { 7 break; 8 } 9}

这个RetentionTimeInMSInDeltaQueue默认是180s,可以看出这个队列是一个长度为180s的滑动窗口,保存最近180s以内的应用实例信息修改,后面我们会看到,客户端调用获取增量信息,实际上就是从这个queue中读取,所以可能一段时间内读取到的信息都是一样的。

可以看出registry使我们所有应用所有实例注册信息保存的地方;但是在客户端从EurekaServer获取实例信息的时候,并不是直接读取registry,而是从Response缓存中获取。

Response缓存的实现类是ResponseCacheImpl,主要包括如下缓存field:

1private final ConcurrentMap<Key, Value> readOnlyCacheMap = new ConcurrentHashMap<Key, Value>(); 2private final LoadingCache<Key, Value> readWriteCacheMap;

一个是guava的loadingcache,一个是普通的ConcurrentHashMap

这个loadingcache的初始化:

1this.readWriteCacheMap = CacheBuilder.newBuilder().initialCapacity(1000) 2 .expireAfterWrite(serverConfig.getResponseCacheAutoExpirationInSeconds(), TimeUnit.SECONDS) 3 .removalListener(new RemovalListener<Key, Value>() { 4 @Override 5 public void onRemoval(RemovalNotification<Key, Value> notification) { 6 Key removedKey = notification.getKey(); 7 if (removedKey.hasRegions()) { 8 Key cloneWithNoRegions = removedKey.cloneWithoutRegions(); 9 regionSpecificKeys.remove(cloneWithNoRegions, removedKey); 10 } 11 } 12 }) 13 .build(new CacheLoader<Key, Value>() { 14 @Override 15 public Value load(Key key) throws Exception { 16 if (key.hasRegions()) { 17 Key cloneWithNoRegions = key.cloneWithoutRegions(); 18 regionSpecificKeys.put(cloneWithNoRegions, key); 19 } 20 Value value = generatePayload(key); 21 return value; 22 } 23 });

对于每个不存在的Key,会首先初始化,主要是调用generatePayload这个方法:

1private Value generatePayload(Key key) { 2 Stopwatch tracer = null; 3 try { 4 String payload; 5 switch (key.getEntityType()) { 6 case Application: 7 boolean isRemoteRegionRequested = key.hasRegions(); 8 9 if (ALL_APPS.equals(key.getName())) { 10 //获取所有应用信息 11 if (isRemoteRegionRequested) { 12 tracer = serializeAllAppsWithRemoteRegionTimer.start(); 13 payload = getPayLoad(key, registry.getApplicationsFromMultipleRegions(key.getRegions())); 14 } else { 15 tracer = serializeAllAppsTimer.start(); 16 payload = getPayLoad(key, registry.getApplications()); 17 } 18 } else if (ALL_APPS_DELTA.equals(key.getName())) { 19 //获取所有应用增量信息 20 if (isRemoteRegionRequested) { 21 tracer = serializeDeltaAppsWithRemoteRegionTimer.start(); 22 versionDeltaWithRegions.incrementAndGet(); 23 versionDeltaWithRegionsLegacy.incrementAndGet(); 24 payload = getPayLoad(key, 25 registry.getApplicationDeltasFromMultipleRegions(key.getRegions())); 26 } else { 27 tracer = serializeDeltaAppsTimer.start(); 28 versionDelta.incrementAndGet(); 29 versionDeltaLegacy.incrementAndGet(); 30 payload = getPayLoad(key, registry.getApplicationDeltas()); 31 } 32 } else { 33 //获取单个应用信息 34 tracer = serializeOneApptimer.start(); 35 payload = getPayLoad(key, registry.getApplication(key.getName())); 36 } 37 break; 38 39 //其他类型我们不关心,先忽略掉相关代码 40 } 41 return new Value(payload); 42 } finally { 43 if (tracer != null) { 44 tracer.stop(); 45 } 46 } 47}

获取所有应用信息,是从registry中直接拿registry.getApplications(),核心方法是getApplicationsFromMultipleRegions,看下简化过的源码:

1public Applications getApplicationsFromMultipleRegions(String[] remoteRegions) { 2 3 boolean includeRemoteRegion = null != remoteRegions && remoteRegions.length != 0; 4 5 Applications apps = new Applications(); 6 apps.setVersion(1L); 7 //将registry中的信息封装好放入Applications 8 for (Entry<String, Map<String, Lease<InstanceInfo>>> entry : registry.entrySet()) { 9 Application app = null; 10 11 if (entry.getValue() != null) { 12 for (Entry<String, Lease<InstanceInfo>> stringLeaseEntry : entry.getValue().entrySet()) { 13 Lease<InstanceInfo> lease = stringLeaseEntry.getValue(); 14 if (app == null) { 15 app = new Application(lease.getHolder().getAppName()); 16 } 17 app.addInstance(decorateInstanceInfo(lease)); 18 } 19 } 20 if (app != null) { 21 apps.addApplication(app); 22 } 23 } 24 //读取其他Region的Apps信息,我们目前不关心,略过这部分代码...... 25 26 //设置AppsHashCode,在之后的介绍中,我们会提到,客户端读取到之后会对比这个AppsHashCode 27 apps.setAppsHashCode(apps.getReconcileHashCode()); 28 return apps; 29}

获取所有应用增量信息,registry.getApplicationDeltas():

1public Applications getApplicationDeltas() { 2 Applications apps = new Applications(); 3 apps.setVersion(responseCache.getVersionDelta().get()); 4 Map<String, Application> applicationInstancesMap = new HashMap<String, Application>(); 5 try { 6 //这里读取用的是写锁,下面我们就会解释为何这么用 7 write.lock(); 8 9 //遍历recentlyChangedQueue,获取所有增量信息 10 Iterator<RecentlyChangedItem> iter = this.recentlyChangedQueue.iterator(); 11 logger.debug("The number of elements in the delta queue is :" 12 + this.recentlyChangedQueue.size()); 13 while (iter.hasNext()) { 14 Lease<InstanceInfo> lease = iter.next().getLeaseInfo(); 15 InstanceInfo instanceInfo = lease.getHolder(); 16 Object[] args = {instanceInfo.getId(), 17 instanceInfo.getStatus().name(), 18 instanceInfo.getActionType().name()}; 19 logger.debug( 20 "The instance id %s is found with status %s and actiontype %s", 21 args); 22 Application app = applicationInstancesMap.get(instanceInfo 23 .getAppName()); 24 if (app == null) { 25 app = new Application(instanceInfo.getAppName()); 26 applicationInstancesMap.put(instanceInfo.getAppName(), app); 27 apps.addApplication(app); 28 } 29 app.addInstance(decorateInstanceInfo(lease)); 30 } 31 32 //读取其他Region的Apps信息,我们目前不关心,略过这部分代码...... 33 34 Applications allApps = getApplications(!disableTransparentFallback); 35 //设置AppsHashCode,在之后的介绍中,我们会提到,客户端读取到之后更新好自己的Apps缓存之后会对比这个AppsHashCode,如果不一样,就会进行一次全量Apps信息请求 36 apps.setAppsHashCode(allApps.getReconcileHashCode()); 37 return apps; 38 } finally { 39 write.unlock(); 40 } 41}

为何这里读写锁这么用,首先我们来分析下这个锁保护的对象是谁,可以很明显的看出,是recentlyChangedQueue这个队列。那么谁在修改这个队列,谁又在读取呢?
每个服务实例注册,取消的时候,都会修改这个队列,这个队列是多线程修改的。但是读取,只有loadingcache的ALL_APPS_DELTAkey初始化线程会读取,而且在缓存失效前都不会再有线程读取。所以可以归纳为,多线程频繁修改,但是单线程不频繁读取。
如果没有锁,那么recentlyChangedQueue在遍历读取时如果遇到修改,就会抛出并发修改异常。如果用writeLock锁住多线程修改,那么同一时间只有一个线程能修改,效率不好。所以。利用读锁锁住多线程修改,利用写锁锁住单线程读取正好符合这里的场景。

前面提到,EurekaClient的查询请求,都是从ResponseCache中获取(从ResponseCache本身缓存的就是请求)。ResponseCache还包括readOnlyCacheMap,这个默认时启用的,就是用户请求会先从readOnlyCacheMap读取,如果readOnlyCacheMap中不存在,则从上面介绍的readWriteCacheMap中获取,之后再放入readOnlyCacheMap。

1Value getValue(final Key key, boolean useReadOnlyCache) { 2 Value payload = null; 3 try { 4 if (useReadOnlyCache) { 5 final Value currentPayload = readOnlyCacheMap.get(key); 6 if (currentPayload != null) { 7 payload = currentPayload; 8 } else { 9 payload = readWriteCacheMap.get(key); 10 readOnlyCacheMap.put(key, payload); 11 } 12 } else { 13 payload = readWriteCacheMap.get(key); 14 } 15 } catch (Throwable t) { 16 logger.error("Cannot get value for key :" + key, t); 17 } 18 return payload; 19}

还有个定时任务:每隔只读缓存刷新时间将ReadWriteMap的信息复制到ReadOnlyMap上面:这个readOnlyCacheMap里面数据是定时从readWriteCacheMap中拷贝出来的:

1 private TimerTask getCacheUpdateTask() { 2 return new TimerTask() { 3 @Override 4 public void run() { 5 logger.debug("Updating the client cache from response cache"); 6 for (Key key : readOnlyCacheMap.keySet()) { 7 if (logger.isDebugEnabled()) { 8 Object[] args = {key.getEntityType(), key.getName(), key.getVersion(), key.getType()}; 9 logger.debug("Updating the client cache from response cache for key : {} {} {} {}", args); 10 } 11 try { 12 CurrentRequestVersion.set(key.getVersion()); 13 Value cacheValue = readWriteCacheMap.get(key); 14 Value currentCacheValue = readOnlyCacheMap.get(key); 15 if (cacheValue != currentCacheValue) { 16 readOnlyCacheMap.put(key, cacheValue); 17 } 18 } catch (Throwable th) { 19 logger.error("Error while updating the client cache from response cache", th); 20 } 21 } 22 } 23 }; 24}

在本篇最开始的时候提到register和cancel都会主动失效对应的ResponseCache,这个主动失效的源代码是:

1public void invalidate(String appName, @Nullable String vipAddress, @Nullable String secureVipAddress) { 2 for (Key.KeyType type : Key.KeyType.values()) { for (Version v : Version.values()) { //对于任意一个APP缓存失效,都要让对应的APP请求响应,全量APP信息请求响应,增量APP信息请求响应失效 invalidate( new Key(Key.EntityType.Application, appName, type, v, EurekaAccept.full), new Key(Key.EntityType.Application, appName, type, v, EurekaAccept.compact), new Key(Key.EntityType.Application, ALL_APPS, type, v, EurekaAccept.full), new Key(Key.EntityType.Application, ALL_APPS, type, v, EurekaAccept.compact), new Key(Key.EntityType.Application, ALL_APPS_DELTA, type, v, EurekaAccept.full), new Key(Key.EntityType.Application, ALL_APPS_DELTA, type, v, EurekaAccept.compact) ); if (null != vipAddress) { invalidate(new Key(Key.EntityType.VIP, vipAddress, type, v, EurekaAccept.full)); } 3 if (null != secureVipAddress) { 4 invalidate(new Key(Key.EntityType.SVIP, secureVipAddress, type, v, EurekaAccept.full)); 5 } 6 } 7 } 8} 9 10public void invalidate(Key... keys) { 11 for (Key key : keys) { 12 logger.debug("Invalidating the response cache key : {} {} {} {}, {}", 13 key.getEntityType(), key.getName(), key.getVersion(), key.getType(), key.getEurekaAccept()); 14 15 readWriteCacheMap.invalidate(key); 16 Collection<Key> keysWithRegions = regionSpecificKeys.get(key); 17 if (null != keysWithRegions && !keysWithRegions.isEmpty()) { 18 for (Key keysWithRegion : keysWithRegions) { 19 logger.debug("Invalidating the response cache key : {} {} {} {} {}", 20 key.getEntityType(), key.getName(), key.getVersion(), key.getType(), key.getEurekaAccept()); 21 readWriteCacheMap.invalidate(keysWithRegion); 22 } 23 } 24 } 25}

在readWriteCacheMap中使对应的APP请求响应,全量APP信息请求响应,增量APP信息请求响应失效后,下次请求,就会再读取registry生成。对于registry,新加入的应用或者实例会被读取到。对于cancel,退出的应用或者实例也会被去除掉

所以,总结起来,用下面这张图展示下EurekaServer 重要缓存和对应的请求:

image

点赞
收藏

评论区

加载中...

相关推荐

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(

皕杰报表之UUID

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

手写Java HashMap源码

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

Spring Cloud Eureka源代码解析(1)Eureka启动,原生启动与SpringCloudEureka启动异同

Eureka作为服务注册中心对整个微服务架构起着最核心的整合作用,因此对Eureka还是有很大的必要进行深入研究。Eureka1.x版本是纯基于servlet的应用。为了与springcloud结合使用,除了本身eureka代码,还有个粘合模块springcloudnetflixeurekaserver。在我们启动EurekaServer实例

springcloud知识点笔记

Eureka的自我保护机制:默认情况下EurekaClient定时向EurekaServer端发送心跳包,如果EurekaServer在一定时间内没有收到EurekaClient发送的心跳包,便会直接从服务注册列表中剔除该服务(默认90S),但是在短时间丢失大量的服务实例心跳,这时候EurekaServer会开启自我保护机制,不会剔除该服务。Ribbon

Spring Cloud Eureka源代码解析(2) EurekaServer 重要缓存解析 - HelloWorld