mirror of
https://github.com/dromara/hertzbeat.git
synced 2026-09-17 09:40:58 +00:00
[refactor] refactor connect common cache (#2469)
Signed-off-by: tomsun28 <tomsun28@outlook.com> Co-authored-by: crossoverJie <crossoverJie@gmail.com>
This commit is contained in:
+14
-60
@@ -21,11 +21,9 @@ import com.google.common.util.concurrent.ThreadFactoryBuilder;
|
||||
import com.googlecode.concurrentlinkedhashmap.ConcurrentLinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ArrayBlockingQueue;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ScheduledThreadPoolExecutor;
|
||||
import java.util.concurrent.ThreadFactory;
|
||||
import java.util.concurrent.ThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
@@ -34,16 +32,11 @@ import lombok.extern.slf4j.Slf4j;
|
||||
*/
|
||||
@Slf4j
|
||||
public class ConnectionCommonCache<T, C extends AbstractConnection<?>> {
|
||||
|
||||
|
||||
/**
|
||||
* default cache time 200s
|
||||
* default cache time 600s
|
||||
*/
|
||||
private static final long DEFAULT_CACHE_TIMEOUT = 200 * 1000L;
|
||||
|
||||
/**
|
||||
* default max cache num
|
||||
*/
|
||||
private static final int DEFAULT_MAX_CAPACITY = 10000;
|
||||
private static final long DEFAULT_CACHE_TIMEOUT = 600 * 1000L;
|
||||
|
||||
/**
|
||||
* cacheTime length
|
||||
@@ -60,19 +53,14 @@ public class ConnectionCommonCache<T, C extends AbstractConnection<?>> {
|
||||
*/
|
||||
private ConcurrentLinkedHashMap<T, C> cacheMap;
|
||||
|
||||
/**
|
||||
* the executor who clean cache when timeout
|
||||
*/
|
||||
private ThreadPoolExecutor timeoutCleanerExecutor;
|
||||
|
||||
public ConnectionCommonCache() {
|
||||
init();
|
||||
initCache();
|
||||
}
|
||||
|
||||
private void init() {
|
||||
private void initCache() {
|
||||
cacheMap = new ConcurrentLinkedHashMap
|
||||
.Builder<T, C>()
|
||||
.maximumWeightedCapacity(DEFAULT_MAX_CAPACITY)
|
||||
.maximumWeightedCapacity(Integer.MAX_VALUE)
|
||||
.listener((key, value) -> {
|
||||
timeoutMap.remove(key);
|
||||
try {
|
||||
@@ -82,70 +70,37 @@ public class ConnectionCommonCache<T, C extends AbstractConnection<?>> {
|
||||
}
|
||||
log.info("connection common cache discard key: {}, value: {}.", key, value);
|
||||
}).build();
|
||||
timeoutMap = new ConcurrentHashMap<>(DEFAULT_MAX_CAPACITY >> 6);
|
||||
// last-first-coverage algorithm, run the first and last thread, discard mid
|
||||
timeoutCleanerExecutor = new ThreadPoolExecutor(1, 1, 1, TimeUnit.SECONDS,
|
||||
new ArrayBlockingQueue<>(1),
|
||||
r -> new Thread(r, "connection-cache-timeout-cleaner"),
|
||||
new ThreadPoolExecutor.DiscardOldestPolicy());
|
||||
timeoutMap = new ConcurrentHashMap<>(16);
|
||||
// init monitor available detector cyc task
|
||||
ThreadFactory threadFactory = new ThreadFactoryBuilder()
|
||||
.setNameFormat("connection-cache-ava-detector-%d")
|
||||
.setNameFormat("connection-cache-timout-detector-%d")
|
||||
.setDaemon(true)
|
||||
.build();
|
||||
ScheduledThreadPoolExecutor scheduledExecutor = new ScheduledThreadPoolExecutor(1, threadFactory);
|
||||
scheduledExecutor.scheduleWithFixedDelay(this::detectCacheAvailable, 2, 20, TimeUnit.MINUTES);
|
||||
scheduledExecutor.scheduleWithFixedDelay(this::cleanTimeoutOrUnHealthCache, 2, 100, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
/**
|
||||
* detect all cache available, cleanup not ava connection
|
||||
* clean and remove timeout cache
|
||||
*/
|
||||
private void detectCacheAvailable() {
|
||||
try {
|
||||
cacheMap.forEach((key, value) -> {
|
||||
Long[] cacheTime = timeoutMap.get(key);
|
||||
long currentTime = System.currentTimeMillis();
|
||||
if (cacheTime == null || cacheTime.length != CACHE_TIME_LENGTH
|
||||
|| cacheTime[0] + cacheTime[1] < currentTime) {
|
||||
cacheMap.remove(key);
|
||||
timeoutMap.remove(key);
|
||||
try {
|
||||
value.close();
|
||||
} catch (Exception e) {
|
||||
log.error("connection close error: {}.", e.getMessage(), e);
|
||||
}
|
||||
|
||||
}
|
||||
});
|
||||
} catch (Exception e) {
|
||||
log.error("connection common cache detect cache available error: {}.", e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* clean timeout cache
|
||||
*/
|
||||
private void cleanTimeoutCache() {
|
||||
private void cleanTimeoutOrUnHealthCache() {
|
||||
try {
|
||||
cacheMap.forEach((key, value) -> {
|
||||
// index 0 is startTime, 1 is timeDiff
|
||||
Long[] cacheTime = timeoutMap.get(key);
|
||||
long currentTime = System.currentTimeMillis();
|
||||
if (cacheTime == null || cacheTime.length != CACHE_TIME_LENGTH) {
|
||||
timeoutMap.put(key, new Long[]{currentTime, DEFAULT_CACHE_TIMEOUT});
|
||||
} else if (cacheTime[0] + cacheTime[1] < currentTime) {
|
||||
// timeout, remove this object cache
|
||||
if (cacheTime == null || cacheTime.length != CACHE_TIME_LENGTH
|
||||
|| cacheTime[0] + cacheTime[1] < currentTime) {
|
||||
log.warn("[connection common cache] clean the timeout cache, key {}", key);
|
||||
timeoutMap.remove(key);
|
||||
cacheMap.remove(key);
|
||||
try {
|
||||
value.close();
|
||||
} catch (Exception e) {
|
||||
log.error("connection close error: {}.", e.getMessage(), e);
|
||||
log.error("clean connection close error: {}.", e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
});
|
||||
Thread.sleep(20 * 1000);
|
||||
} catch (Exception e) {
|
||||
log.error("[connection common cache] clean timeout cache error: {}.", e.getMessage(), e);
|
||||
}
|
||||
@@ -165,7 +120,6 @@ public class ConnectionCommonCache<T, C extends AbstractConnection<?>> {
|
||||
}
|
||||
cacheMap.put(key, value);
|
||||
timeoutMap.put(key, new Long[]{System.currentTimeMillis(), timeDiff});
|
||||
timeoutCleanerExecutor.execute(this::cleanTimeoutCache);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user