如何实现Guava Cache的非阻塞get()与异步加载?
Guava Cache实现非阻塞get()与异步刷新方案
问题背景
需要基于Guava Cache实现全链路非阻塞的缓存读取:
- 缓存值从数据库异步加载,数据库查询成本高
- 允许在后台刷新新数据时返回过期缓存
- 绝对禁止读取缓存时发生阻塞
最初尝试结合refreshAfterWrite和expireAfterWrite,但发现当缓存条目需要刷新时,cache.get()会阻塞——根据Guava文档:
使用CacheBuilder.refreshAfterWrite(long, TimeUnit)可为缓存添加自动定时刷新。与expireAfterWrite不同,refreshAfterWrite会让键在指定时长后具备刷新资格,但仅在查询条目时才会实际启动刷新。
原实现代码:
class DbCacheLoader() extends CacheLoader[K, V] { override def load(key: K): V = { // read from DB } }
val syncLoader = new DbCacheLoader() val cacheLoader = CacheLoader.asyncReloading(syncLoader, threadPool) val cache: LoadingCache[K, V] = CacheBuilder.newBuilder() .refreshAfterWrite(timeout, TimeUnit.MILLISECONDS) // timeout .expireAfterWrite(hardTimeout, TimeUnit.MILLISECONDS) // hard timeout .concurrencyLevel(threadPoolSize * 10) .recordStats() .build[K, V](cacheLoader)
核心问题分析
原方案的阻塞根源在于使用了LoadingCache而非AsyncLoadingCache:
LoadingCache的get()是同步方法,首次加载缓存时会直接阻塞调用线程执行load()逻辑- 即使使用
CacheLoader.asyncReloading,刷新阶段确实会异步执行,但首次加载的阻塞无法避免
解决方案:使用AsyncLoadingCache实现全非阻塞
Guava提供的AsyncLoadingCache是专门为异步加载场景设计的,它的所有加载/刷新操作都通过ListenableFuture异步执行,调用线程完全不会被阻塞。
1. 修正缓存构建代码
class DbCacheLoader() extends CacheLoader[K, V] { override def load(key: K): V = { // 同步从数据库加载数据的逻辑 fetchDataFromDb(key) } private def fetchDataFromDb(key: K): V = { // 实际数据库查询实现 } }
val syncLoader = new DbCacheLoader() // 用asyncReloading包装同步加载器,指定后台线程池 val asyncCacheLoader = CacheLoader.asyncReloading(syncLoader, threadPool) // 构建AsyncLoadingCache而非LoadingCache val asyncCache: AsyncLoadingCache[K, V] = CacheBuilder.newBuilder() .refreshAfterWrite(timeout, TimeUnit.MILLISECONDS) // 软刷新间隔:到点后具备刷新资格,查询时异步刷新 .expireAfterWrite(hardTimeout, TimeUnit.MILLISECONDS) // 硬过期时间:过期后直接移除缓存,触发全新异步加载 .concurrencyLevel(threadPoolSize * 10) .recordStats() .buildAsync(asyncCacheLoader)
2. 非阻塞的缓存读取方式
调用AsyncLoadingCache.get()会立即返回ListenableFuture,不会阻塞当前线程,你可以通过回调处理加载结果:
// 非阻塞获取缓存,立即返回Future val futureValue: ListenableFuture[V] = asyncCache.get(key) // 添加异步回调处理结果 Futures.addCallback(futureValue, new FutureCallback[V] { override def onSuccess(result: V): Unit = { // 处理加载成功的缓存值 } override def onFailure(t: Throwable): Unit = { // 处理数据库加载失败的异常 } }, threadPool)
3. 关键行为说明
- 首次加载:缓存中无对应条目时,
get()会触发异步加载,返回Future,调用线程无需等待,完全非阻塞 - 后台刷新:当条目过了
refreshAfterWrite时间,下一次get()会立即返回旧缓存值,同时后台线程异步执行新数据加载,加载完成后自动更新缓存 - 硬过期处理:条目过了
expireAfterWrite时间后,缓存会被移除,此时get()会触发全新的异步加载,同样返回Future不阻塞
内容的提问来源于stack exchange,提问作者bachr
相关产品推荐
相关产品推荐

