You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 04:45:35