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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 02:56:52