diff --git a/orm/src/main/java/com/zfoo/orm/cache/EntityCache.java b/orm/src/main/java/com/zfoo/orm/cache/EntityCache.java index 0dfdded3..32092258 100644 --- a/orm/src/main/java/com/zfoo/orm/cache/EntityCache.java +++ b/orm/src/main/java/com/zfoo/orm/cache/EntityCache.java @@ -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, E extends IEntity> imple private final Class clazz; private final EntityDef entityDef; - private final LazyCache> cache; + private final LazyCache> caches; private final IEntityWrapper wrapper; @@ -69,17 +68,17 @@ public class EntityCache, E extends IEntity> imple wrapper = EnhanceUtils.createEntityWrapper(entityWrapper); } - var removeCallback = new BiConsumer>>, LazyCache.RemovalCause>() { + var removeCallback = new BiConsumer>>, LazyCache.RemovalCause>() { @Override - public void accept(List>> removes, LazyCache.RemovalCause removalCause) { + public void accept(List>> 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, E extends IEntity> 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, E extends IEntity> 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, E extends IEntity> 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, E extends IEntity> imple */ private PNode 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, E extends IEntity> 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, E extends IEntity> imple var currentTime = TimeUtils.currentTimeMillis(); // key为threadId var updateMap = new HashMap>(); - var initSize = cache.size() >> 2; - cache.forEach(new BiConsumer>() { + var initSize = caches.size() >> 2; + caches.forEach(new BiConsumer>() { @Override public void accept(PK pk, PNode pnode) { var entity = pnode.getEntity(); @@ -267,8 +266,8 @@ public class EntityCache, E extends IEntity> imple @Override public void persistAllBlock() { var currentTime = TimeUtils.currentTimeMillis(); - var updateList = new ArrayList(cache.size()); - cache.forEach(new BiConsumer>() { + var updateList = new ArrayList(caches.size()); + caches.forEach(new BiConsumer>() { @Override public void accept(PK pk, PNode pnode) { var entity = pnode.getEntity(); @@ -355,7 +354,7 @@ public class EntityCache, E extends IEntity> 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, E extends IEntity> 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, E extends IEntity> imple @Override public void forEach(BiConsumer 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(); } } diff --git a/scheduler/src/main/java/com/zfoo/scheduler/util/LazyCache.java b/scheduler/src/main/java/com/zfoo/scheduler/util/LazyCache.java index 9201d7f4..41cf270b 100644 --- a/scheduler/src/main/java/com/zfoo/scheduler/util/LazyCache.java +++ b/scheduler/src/main/java/com/zfoo/scheduler/util/LazyCache.java @@ -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 { private static final float DEFAULT_BACK_PRESSURE_FACTOR = 0.13f; - private static class CacheValue { - public volatile V value; + public static class Cache { + public K k; + public V v; public volatile long expireTime; } @@ -59,11 +59,11 @@ public class LazyCache { private long expireCheckIntervalMillis; private volatile long minExpireTime; private AtomicLong expireCheckTimeAtomic; - private ConcurrentMap> cacheMap; - private BiConsumer>, RemovalCause> removeListener = (removes, removalCause) -> { + private ConcurrentMap> cacheMap; + private BiConsumer>, RemovalCause> removeListener = (removes, removalCause) -> { }; - public LazyCache(int maximumSize, long expireAfterAccessMillis, long expireCheckIntervalMillis, BiConsumer>, RemovalCause> removeListener) { + public LazyCache(int maximumSize, long expireAfterAccessMillis, long expireCheckIntervalMillis, BiConsumer>, RemovalCause> removeListener) { AssertionUtils.ge1(maximumSize); AssertionUtils.ge0(expireAfterAccessMillis); AssertionUtils.ge0(expireCheckIntervalMillis); @@ -83,30 +83,32 @@ public class LazyCache { * 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(); - 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(); + 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 { public void forEach(BiConsumer 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 { 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> list, RemovalCause removalCause) { + private void removeForCause(List> 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 { if (expireCheckTimeAtomic.compareAndSet(expireCheckTime, now + expireCheckIntervalMillis)) { if (now > this.minExpireTime) { var minTimestamp = Long.MAX_VALUE; - var removeList = new ArrayList>(); - for (var entry : cacheMap.entrySet()) { - var expireTime = entry.getValue().expireTime; + var removeList = new ArrayList>(); + 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) { diff --git a/scheduler/src/test/java/com/zfoo/scheduler/util/LazyCacheTesting.java b/scheduler/src/test/java/com/zfoo/scheduler/util/LazyCacheTesting.java index 4c72fe3d..795de78a 100644 --- a/scheduler/src/test/java/com/zfoo/scheduler/util/LazyCacheTesting.java +++ b/scheduler/src/test/java/com/zfoo/scheduler/util/LazyCacheTesting.java @@ -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>, LazyCache.RemovalCause> myRemoveCallback = new BiConsumer>, LazyCache.RemovalCause>() { + private static final BiConsumer>, LazyCache.RemovalCause> myRemoveCallback = new BiConsumer>, LazyCache.RemovalCause>() { @Override - public void accept(List> pairs, LazyCache.RemovalCause removalCause) { - for (var pair : pairs) { - logger.info("remove key:[{}] value:[{}] removalCause:[{}]", pair.getKey(), pair.getValue(), removalCause); + public void accept(List> caches, LazyCache.RemovalCause removalCause) { + for (var cache : caches) { + logger.info("remove key:[{}] value:[{}] removalCause:[{}]", cache.k, cache.v, removalCause); } } };