perf[cache]: reduce the duplication of objects

This commit is contained in:
godotg committed 2024-08-20 13:27:28 +08:00
1 parent d676bc4972
commit 1277ab03ea
3 files changed
+58 -59

No files matched your search

+21 -22
View File
@@ -28,7 +28,6 @@ import com.zfoo.orm.model.IEntity;
import com.zfoo.orm.query.Page;
import com.zfoo.protocol.collection.CollectionUtils;
import com.zfoo.protocol.exception.RunException;
import com.zfoo.protocol.model.Pair;
import com.zfoo.protocol.util.*;
import com.zfoo.scheduler.manager.SchedulerBus;
import com.zfoo.scheduler.util.LazyCache;
@@ -54,7 +53,7 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
private final Class<E> clazz;
private final EntityDef entityDef;
private final LazyCache<PK, PNode<PK, E>> cache;
private final LazyCache<PK, PNode<PK, E>> caches;
private final IEntityWrapper<PK, E> wrapper;
@@ -69,17 +68,17 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
wrapper = EnhanceUtils.createEntityWrapper(entityWrapper);
}
var removeCallback = new BiConsumer<List<Pair<PK, PNode<PK, E>>>, LazyCache.RemovalCause>() {
var removeCallback = new BiConsumer<List<LazyCache.Cache<PK, PNode<PK, E>>>, LazyCache.RemovalCause>() {
@Override
public void accept(List<Pair<PK, PNode<PK, E>>> removes, LazyCache.RemovalCause removalCause) {
public void accept(List<LazyCache.Cache<PK, PNode<PK, E>>> removes, LazyCache.RemovalCause removalCause) {
if (removalCause == LazyCache.RemovalCause.EXPLICIT) {
return;
}
EventBus.asyncExecute(clazz.hashCode(), () -> doPersist(removes.stream().map(it -> it.getValue().getEntity()).toList()));
EventBus.asyncExecute(clazz.hashCode(), () -> doPersist(removes.stream().map(it -> it.v.getEntity()).toList()));
}
};
var expireCheckIntervalMillis = Math.max(3 * TimeUtils.MILLIS_PER_SECOND, entityDef.getExpireMillisecond() / 10);
this.cache = new LazyCache<>(entityDef.getCacheSize(), entityDef.getExpireMillisecond(), expireCheckIntervalMillis, removeCallback);
this.caches = new LazyCache<>(entityDef.getCacheSize(), entityDef.getExpireMillisecond(), expireCheckIntervalMillis, removeCallback);
if (CollectionUtils.isNotEmpty(entityDef.getIndexDefMap())) {
// indexMap
@@ -98,7 +97,7 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
@Override
public E load(PK pk) {
AssertionUtils.notNull(pk);
var pnode = cache.get(pk);
var pnode = caches.get(pk);
if (pnode != null) {
return pnode.getEntity();
}
@@ -110,14 +109,14 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
logger.warn("[{}] can not load [pk:{}] and use null to replace it", clazz.getSimpleName(), pk);
}
pnode = new PNode<>(entity);
cache.put(pk, pnode);
caches.put(pk, pnode);
return entity;
}
@Override
public E loadOrCreate(PK pk) {
AssertionUtils.notNull(pk);
var pnode = cache.get(pk);
var pnode = caches.get(pk);
if (pnode != null) {
return pnode.getEntity();
}
@@ -130,7 +129,7 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
OrmContext.getAccessor().insert(entity);
}
pnode = new PNode<>(entity);
cache.put(pk, pnode);
caches.put(pk, pnode);
return entity;
}
@@ -139,10 +138,10 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
*/
private PNode<PK, E> fetchCachePnode(E entity, boolean safe) {
var id = entity.id();
var cachePnode = cache.get(id);
var cachePnode = caches.get(id);
if (cachePnode == null) {
cachePnode = new PNode<>(entity);
cache.put(entity.id(), cachePnode);
caches.put(entity.id(), cachePnode);
}
// 比较地址是否相等
@@ -209,12 +208,12 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
// 游戏业务中,操作最频繁的是update,不是insert,delete,query
// 所以这边并不考虑
AssertionUtils.notNull(pk);
cache.remove(pk);
caches.remove(pk);
}
@Override
public void persist(PK pk) {
var pnode = cache.get(pk);
var pnode = caches.get(pk);
if (pnode == null) {
return;
}
@@ -237,8 +236,8 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
var currentTime = TimeUtils.currentTimeMillis();
// key为threadId
var updateMap = new HashMap<Long, List<E>>();
var initSize = cache.size() >> 2;
cache.forEach(new BiConsumer<PK, PNode<PK, E>>() {
var initSize = caches.size() >> 2;
caches.forEach(new BiConsumer<PK, PNode<PK, E>>() {
@Override
public void accept(PK pk, PNode<PK, E> pnode) {
var entity = pnode.getEntity();
@@ -267,8 +266,8 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
@Override
public void persistAllBlock() {
var currentTime = TimeUtils.currentTimeMillis();
var updateList = new ArrayList<E>(cache.size());
cache.forEach(new BiConsumer<PK, PNode<PK, E>>() {
var updateList = new ArrayList<E>(caches.size());
caches.forEach(new BiConsumer<PK, PNode<PK, E>>() {
@Override
public void accept(PK pk, PNode<PK, E> pnode) {
var entity = pnode.getEntity();
@@ -355,7 +354,7 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
var dbEntity = dbMap.get(id);
if (dbEntity == null) {
cache.remove(entity.id());
caches.remove(entity.id());
logger.warn("[database:{}] not found entity [id:{}]", clazz.getSimpleName(), id);
continue;
}
@@ -376,7 +375,7 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
}
// 数据库版本号较大,说明缓存的数据不是最新的,直接清除缓存,下次重新加载
cache.remove(id);
caches.remove(id);
load(id);
logger.warn("[database:{}] document of entity [id:{}] version [{}] is greater than cache [vs:{}] and reload db entity to cache", clazz.getSimpleName(), id, dbEntityVersion, entityVersion);
}
@@ -384,12 +383,12 @@ public class EntityCache<PK extends Comparable<PK>, E extends IEntity<PK>> imple
@Override
public void forEach(BiConsumer<PK, E> biConsumer) {
cache.forEach((pk, pnode) -> biConsumer.accept(pk, pnode.getEntity()));
caches.forEach((pk, pnode) -> biConsumer.accept(pk, pnode.getEntity()));
}
@Override
public long size() {
return cache.size();
return caches.size();
}
}
@@ -1,7 +1,6 @@
package com.zfoo.scheduler.util;
import com.zfoo.protocol.collection.CollectionUtils;
import com.zfoo.protocol.model.Pair;
import com.zfoo.protocol.util.AssertionUtils;
import java.util.ArrayList;
@@ -20,8 +19,9 @@ public class LazyCache<K, V> {
private static final float DEFAULT_BACK_PRESSURE_FACTOR = 0.13f;
private static class CacheValue<V> {
public volatile V value;
public static class Cache<K, V> {
public K k;
public V v;
public volatile long expireTime;
}
@@ -59,11 +59,11 @@ public class LazyCache<K, V> {
private long expireCheckIntervalMillis;
private volatile long minExpireTime;
private AtomicLong expireCheckTimeAtomic;
private ConcurrentMap<K, CacheValue<V>> cacheMap;
private BiConsumer<List<Pair<K, V>>, RemovalCause> removeListener = (removes, removalCause) -> {
private ConcurrentMap<K, Cache<K, V>> cacheMap;
private BiConsumer<List<Cache<K, V>>, RemovalCause> removeListener = (removes, removalCause) -> {
};
public LazyCache(int maximumSize, long expireAfterAccessMillis, long expireCheckIntervalMillis, BiConsumer<List<Pair<K, V>>, RemovalCause> removeListener) {
public LazyCache(int maximumSize, long expireAfterAccessMillis, long expireCheckIntervalMillis, BiConsumer<List<Cache<K, V>>, RemovalCause> removeListener) {
AssertionUtils.ge1(maximumSize);
AssertionUtils.ge0(expireAfterAccessMillis);
AssertionUtils.ge0(expireCheckIntervalMillis);
@@ -83,30 +83,32 @@ public class LazyCache<K, V> {
* If the cache previously contained a value associated with the key, the old value is replaced by the new value.
*/
public void put(K key, V value) {
var cacheValue = new CacheValue<V>();
cacheValue.value = value;
cacheValue.expireTime = TimeUtils.now() + expireAfterAccessMillis;
var oldCacheValue = cacheMap.put(key, cacheValue);
if (oldCacheValue != null) {
removeListener.accept(List.of(new Pair<>(key, oldCacheValue.value)), RemovalCause.REPLACED);
}
checkExpire();
checkMaximumSize();
var cache = new Cache<K, V>();
cache.k = key;
cache.v = value;
cache.expireTime = TimeUtils.now() + expireAfterAccessMillis;
var oldCache = cacheMap.put(key, cache);
if (oldCache != null) {
removeListener.accept(List.of(oldCache), RemovalCause.REPLACED);
}
}
public V get(K key) {
checkExpire();
var cacheValue = cacheMap.get(key);
if (cacheValue == null) {
var cache = cacheMap.get(key);
if (cache == null) {
return null;
}
if (cacheValue.expireTime < TimeUtils.now()) {
if (cache.expireTime < TimeUtils.now()) {
removeForCause(key, RemovalCause.EXPIRED);
return null;
}
cacheValue.expireTime = TimeUtils.now() + expireAfterAccessMillis;
return cacheValue.value;
cache.expireTime = TimeUtils.now() + expireAfterAccessMillis;
return cache.v;
}
@@ -117,8 +119,8 @@ public class LazyCache<K, V> {
public void forEach(BiConsumer<K, V> biConsumer) {
checkExpire();
for (var entry : cacheMap.entrySet()) {
biConsumer.accept(entry.getKey(), entry.getValue().value);
for (var cache : cacheMap.values()) {
biConsumer.accept(cache.k, cache.v);
}
}
@@ -133,29 +135,28 @@ public class LazyCache<K, V> {
if (key == null) {
return;
}
var cacheValue = cacheMap.remove(key);
if (cacheValue != null) {
removeListener.accept(List.of(new Pair<>(key, cacheValue.value)), removalCause);
var cache = cacheMap.remove(key);
if (cache != null) {
removeListener.accept(List.of(cache), removalCause);
}
}
private void removeForCause(List<Pair<K, V>> list, RemovalCause removalCause) {
private void removeForCause(List<Cache<K, V>> list, RemovalCause removalCause) {
if (CollectionUtils.isEmpty(list)) {
return;
}
var removeList = list.stream()
.filter(it -> cacheMap.remove(it.getKey()) != null)
.filter(it -> cacheMap.remove(it.k) != null)
.toList();
removeListener.accept(removeList, removalCause);
}
private void checkMaximumSize() {
if (cacheMap.size() > backPressureSize) {
var removeList = cacheMap.entrySet()
var removeList = cacheMap.values()
.stream()
.sorted((a, b) -> Long.compare(a.getValue().expireTime, b.getValue().expireTime))
.sorted((a, b) -> Long.compare(a.expireTime, b.expireTime))
.limit(Math.max(0, cacheMap.size() - maximumSize))
.map(it -> new Pair<>(it.getKey(), it.getValue().value))
.toList();
removeForCause(removeList, RemovalCause.SIZE);
}
@@ -168,11 +169,11 @@ public class LazyCache<K, V> {
if (expireCheckTimeAtomic.compareAndSet(expireCheckTime, now + expireCheckIntervalMillis)) {
if (now > this.minExpireTime) {
var minTimestamp = Long.MAX_VALUE;
var removeList = new ArrayList<Pair<K, V>>();
for (var entry : cacheMap.entrySet()) {
var expireTime = entry.getValue().expireTime;
var removeList = new ArrayList<Cache<K, V>>();
for (var cache : cacheMap.values()) {
var expireTime = cache.expireTime;
if (expireTime < now) {
removeList.add(new Pair<>(entry.getKey(), entry.getValue().value));
removeList.add(cache);
continue;
}
if (expireTime < minTimestamp) {
@@ -1,6 +1,5 @@
package com.zfoo.scheduler.util;
import com.zfoo.protocol.model.Pair;
import com.zfoo.protocol.util.ThreadUtils;
import org.junit.Ignore;
import org.junit.Test;
@@ -20,11 +19,11 @@ public class LazyCacheTesting {
private static final Logger logger = LoggerFactory.getLogger(LazyCacheTesting.class);
private static final BiConsumer<List<Pair<Integer, String>>, LazyCache.RemovalCause> myRemoveCallback = new BiConsumer<List<Pair<Integer, String>>, LazyCache.RemovalCause>() {
private static final BiConsumer<List<LazyCache.Cache<Integer, String>>, LazyCache.RemovalCause> myRemoveCallback = new BiConsumer<List<LazyCache.Cache<Integer, String>>, LazyCache.RemovalCause>() {
@Override
public void accept(List<Pair<Integer, String>> pairs, LazyCache.RemovalCause removalCause) {
for (var pair : pairs) {
logger.info("remove key:[{}] value:[{}] removalCause:[{}]", pair.getKey(), pair.getValue(), removalCause);
public void accept(List<LazyCache.Cache<Integer, String>> caches, LazyCache.RemovalCause removalCause) {
for (var cache : caches) {
logger.info("remove key:[{}] value:[{}] removalCause:[{}]", cache.k, cache.v, removalCause);
}
}
};