如何在Guava Cache的load方法中异步调用服务并非阻塞加载缓存?
Guava Cache异步非阻塞加载实现方案
我已实现Guava Cache,希望在其CacheLoader的load方法中异步调用服务,且不阻塞其他线程,通过回调完成缓存加载。目前constructDataMangerFromSource方法中通过future.get()阻塞线程,无法实现非阻塞需求,同时希望在异步回调的onError中抛出异常、onSuccess中返回值。
当前实现代码
CacheLoader实现类
public CacheLoader(final DataValueProvider dataValueProvider) { this.parentExecutor = Executors.newSingleThreadExecutor(new DefaultThreadFactory("cache-loader")); this.executorService = MoreExecutors.listeningDecorator(parentExecutor); this.secretValueProvider = secretValueProvider; } /** * Call back method invoked by Guava to loads the initial value of the secretManager entry. */ @Override public String load(CacheIndex cacheIndex) throws Exception { return dataValueProvider.constructDataMangerFromSource(cacheIndex); }
DataValueProvider的constructDataMangerFromSource方法
protected String constructDataMangerFromSource(CacheIndex cacheIndex) { GetSecretValueRequest getSecretValueRequest = new GetSecretValueRequest() .withSecretId(cacheIndex.secretArn); Future<GetSecretValueResult> future = awsSecretsManagerAsync.getSecretValueAsync( getSecretValueRequest, new AsyncHandler<GetSecretValueRequest, GetSecretValueResult>() { @SneakyThrows @Override public void onError(Exception e) { //I want to throw exception from here } @Override public void onSuccess(GetSecretValueRequest request, GetSecretValueResult getSecretValueResult) { long endTime = System.currentTimeMillis(); long totalTime = endTime - startTime; logger.info("Total time taken in the call " + totalTime); //I want to return the value from here } }); try { //This is currently blocking other thread GetSecretValueResult getSecretValueResult = future.get(); return getSecretValueResult.getSecretString(); } catch (Exception ex) { logger.error("SecretManage call failed " + ex.getCause() + " " + ex.getMessage()); throw ex; } }
解决方案:基于Guava ListenableFuture实现异步加载
Guava Cache原生支持异步加载逻辑,核心是让CacheLoader重载返回ListenableFuture<String>的load方法,替代同步返回String的逻辑,从而实现非阻塞的缓存加载。
步骤1:修改CacheLoader的load方法
替换同步load方法为异步版本,利用已有的ListeningExecutorService提交异步任务:
@Override public ListenableFuture<String> load(CacheIndex cacheIndex) throws Exception { // 将加载任务提交到指定线程池,异步执行 return executorService.submit(() -> dataValueProvider.constructDataMangerFromSourceAsync(cacheIndex)); }
步骤2:重构异步加载方法
移除阻塞的future.get(),使用Guava的SettableFuture接收AWS异步调用的结果和异常:
protected ListenableFuture<String> constructDataMangerFromSourceAsync(CacheIndex cacheIndex) { // 创建SettableFuture用于传递异步结果 SettableFuture<String> settableFuture = SettableFuture.create(); GetSecretValueRequest getSecretValueRequest = new GetSecretValueRequest() .withSecretId(cacheIndex.secretArn); awsSecretsManagerAsync.getSecretValueAsync( getSecretValueRequest, new AsyncHandler<GetSecretValueRequest, GetSecretValueResult>() { @Override public void onError(Exception e) { logger.error("SecretManager call failed", e); // 将异常传递给SettableFuture,Guava会在缓存获取时抛出该异常 settableFuture.setException(e); } @Override public void onSuccess(GetSecretValueRequest request, GetSecretValueResult getSecretValueResult) { long totalTime = System.currentTimeMillis() - startTime; logger.info("Total time taken in the call: {}", totalTime); // 将加载结果设置到SettableFuture,完成缓存加载 settableFuture.set(getSecretValueResult.getSecretString()); } }); return settableFuture; }
步骤3:创建支持异步加载的LoadingCache
保持CacheBuilder的配置,传入修改后的CacheLoader即可:
LoadingCache<CacheIndex, String> cache = CacheBuilder.newBuilder() .expireAfterWrite(10, TimeUnit.MINUTES) .build(new YourCacheLoader(dataValueProvider));
关键说明
- 非阻塞特性:通过
SettableFuture衔接AWS异步回调,避免了future.get()的线程阻塞,调用线程会立即返回Future,由Guava在后台监听异步任务完成状态。 - 异常传递:
onError中调用settableFuture.setException(e),Guava会将异常传递给缓存的调用方,当调用方尝试获取缓存值时会抛出该异常。 - 线程池复用:复用已创建的
ListeningExecutorService,确保异步加载任务在指定线程池中执行,避免线程资源混乱。
内容的提问来源于stack exchange,提问作者codemaster001
相关产品推荐
相关产品推荐

