背景
在zookeeper服务端中需要管理session和connection这两类具有超时属性的对象。zookeeper提供了ExpiryQueue来实现通用对象超时管理容器。
实现
我拿connection来说如何管理他的超时问题,不同的连接超时的时间点是不同的

那么zookeeper是如何高效的管理这些不同的超时时间点的呢?
在ExpiryQueue有个expirationInterval属性,zookeeper会使用expirationInterval作为基准把每一个连接的超时时间点归一化为expirationInterval整数倍,归一化的计算方式为
我们拿上一个图为例子,通过归一化处理connection_1,connection_2,connection_3这三个连接的超时时间点可能会变成同一个
假定expirationInterval = 10000,
connection_1_timeout_point = 1599715479084
connection_2_timeout_point = 1599715479184
connection_3_timeout_point = 1599715479384
那么通过归一化计算他们三个的超时时间点都变成了1599715480000
我们现在看下ExpiryQueue两个重要的属性
//记录了每一个对象的归一化后的超时时间点,key是被管理的对象,value是超时时间点
private final ConcurrentHashMap<E, Long> elemMap = new ConcurrentHashMap<E, Long>();
//存储了相同超时时间点的所有对象,key是超时时间点,value是相同超时时间点的对象集合
private final ConcurrentHashMap<Long, Set<E>> expiryMap = new ConcurrentHashMap<Long, Set<E>>();
超时管理线程
有了超时对象管理容器,还需要相应的超时管理线程来监控容器中对象的状态,超时管理线程也会按照expirationInterval为时间间隔单位来运行,它每次运行的时间点记录在ExpiryQueue的nextExpirationTime属性中,nextExpirationTime初始化值由
now是ExpiryQueue创建时的系统时间点
nextExpirationTime之后每次更新的值通过下面公式得到
可以看出nextExpirationTime和归一化后的connection超时时间点是一致的,都是expirationInterval的倍数
超时管理线程处理超时对象
超时管理线程会每隔expirationInterval去ExpiryQueue中获取超时对象,具体实现是ExpiryQueue.poll方法
public Set<E> poll() {
long now = Time.currentElapsedTime();
long expirationTime = nextExpirationTime.get();
if (now < expirationTime) {
return Collections.emptySet();
}
Set<E> set = null;
//更新nextExpirationTime = nextExpirationTime + expirationInterval
long newExpirationTime = expirationTime + expirationInterval;
if (nextExpirationTime.compareAndSet(expirationTime, newExpirationTime)) {
//获取expirationTime超时时间点对应的超时对象集合
set = expiryMap.remove(expirationTime);
}
if (set == null) {
return Collections.emptySet();
}
//返回在本次expirationTime超时时间点超时的对象
return set;
}
每次在超时时间点获取超时对象之后,超时管理线程可以根据超时对象的不同业务特性做不同的业务逻辑
超时对象的更新
超时对象在放入到ExpiryQueue中之后,这些对象会根据自己的特性(比如对于连接对象来说当连接上发生IO事件,那么就需要更新连接超时时间点)更新自己的超时时间点,具体的实现逻辑是在ExpiryQueue.update方法中
//timeout是超时对象的超时时间(或者说存活时长)
public Long update(E elem, int timeout) {
//通过elemMap获取超时对象之前的超时时间点
Long prevExpiryTime = elemMap.get(elem);
long now = Time.currentElapsedTime();
//通过归一化方法获取超时对象新的超时时间点
Long newExpiryTime = roundToNextInterval(now + timeout);
//如果新超时时间点和老超时时间点一样那么不做任何处理
if (newExpiryTime.equals(prevExpiryTime)) {
// No change, so nothing to update
return null;
}
// First add the elem to the new expiry time bucket in expiryMap.
//使用新超时时间点从expiryMap获取所有在新超时时间点超时的对象集合
Set<E> set = expiryMap.get(newExpiryTime);
if (set == null) {
// Construct a ConcurrentHashSet using a ConcurrentHashMap
//如果超时对象集合为空,那么创建一个
set = Collections.newSetFromMap(new ConcurrentHashMap<E, Boolean>());
// Put the new set in the map, but only if another thread
// hasn't beaten us to it
//并发的情况下可能会出现多个线程同时创建相同超时时间点对象集合,所以需要做如下是否存在判断处理
Set<E> existingSet = expiryMap.putIfAbsent(newExpiryTime, set);
if (existingSet != null) {
set = existingSet;
}
}
//把本超时对象加入集合
set.add(elem);
// Map the elem to the new expiry time. If a different previous
// mapping was present, clean up the previous expiry bucket.
//同时更新超时对象在elemMap中新的超时时间点
prevExpiryTime = elemMap.put(elem, newExpiryTime);
if (prevExpiryTime != null && !newExpiryTime.equals(prevExpiryTime)) {
//根据超时对象上一个超时时间点从expiryMap对应的超时对象集合中把本超时对象删除
Set<E> prevSet = expiryMap.get(prevExpiryTime);
if (prevSet != null) {
prevSet.remove(elem);
}
}
return newExpiryTime;
}
上面从源码的角度分析了zookeeper是如何实现超时对象管理的,关于这一块的理解强烈推荐大家看《从 Paxos 到 ZooKeeper:分布式一致性原理与实践》这本书,这本书写的太棒了