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

如何在Caffeine Cache中结合Mono实现异步缓存?

Caffeine Cache结合WebFlux Mono实现异步加载

Caffeine的异步加载能力可以和WebFlux的Mono无缝适配,核心是将两者的异步类型(CompletableFuture与Mono)互相转换,保持响应式非阻塞特性。以下是具体实现方案:

1. 核心思路

Caffeine的AsyncLoadingCache默认返回CompletableFuture,而WebFlux的Mono可以通过Mono.fromFuture()将其包装为响应式类型;反之,若数据源返回Mono,则可以通过mono.toFuture()转换为CompletableFuture交给Caffeine处理。

2. 代码实现示例

2.1 创建异步缓存实例

import com.github.benmanes.caffeine.cache.AsyncLoadingCache;
import com.github.benmanes.caffeine.cache.Caffeine;
import reactor.core.publisher.Mono;

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;

public class ReactiveCaffeineCache {
    // 定义异步加载缓存
    private final AsyncLoadingCache<String, String> asyncCache;

    public ReactiveCaffeineCache() {
        this.asyncCache = Caffeine.newBuilder()
                .expireAfterWrite(5, TimeUnit.MINUTES) // 缓存过期时间
                .maximumSize(100) // 缓存最大容量
                .buildAsync((key, executor) -> {
                    // 调用返回Mono的数据源,转换为CompletableFuture交给Caffeine
                    return fetchDataFromReactiveSource(key).toFuture();
                });
    }

    // 模拟响应式数据源(比如数据库查询、远程调用)
    private Mono<String> fetchDataFromReactiveSource(String key) {
        return Mono.just("Cached value for key: " + key)
                .delayElement(java.time.Duration.ofMillis(100)); // 模拟IO延迟
    }
}

2.2 封装响应式查询方法

为外部提供返回Mono的查询接口,将Caffeine的CompletableFuture转换为响应式类型:

// 对外暴露的查询方法,返回Mono
public Mono<String> getCachedValue(String key) {
    return Mono.fromFuture(asyncCache.get(key));
}

2.3 手动更新缓存

如果需要手动向缓存写入数据,同样适配响应式流程:

// 手动写入缓存,返回Mono<Void>表示操作完成
public Mono<Void> putCachedValue(String key, String value) {
    return Mono.fromRunnable(() -> 
        asyncCache.put(key, CompletableFuture.completedFuture(value))
    ).then();
}

// 或者从响应式数据源获取后自动写入缓存
public Mono<String> getAndCacheFromSource(String key) {
    return fetchDataFromReactiveSource(key)
            .doOnSuccess(value -> 
                asyncCache.put(key, CompletableFuture.completedFuture(value))
            );
}

3. 线程池优化建议

Caffeine默认使用ForkJoinPool.commonPool()处理异步加载,为避免与WebFlux的线程池冲突,可自定义线程池:

this.asyncCache = Caffeine.newBuilder()
        .executor(Executors.newFixedThreadPool(8)) // 自定义线程池
        .expireAfterWrite(5, TimeUnit.MINUTES)
        .maximumSize(100)
        .buildAsync((key, executor) -> fetchDataFromReactiveSource(key).toFuture());

也可以结合WebFlux的调度器,让数据源操作在指定线程池执行:

.buildAsync((key, executor) -> 
    fetchDataFromReactiveSource(key)
        .subscribeOn(reactor.core.scheduler.Schedulers.boundedElastic())
        .toFuture()
);

内容的提问来源于stack exchange,提问作者A Jar of Clay

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 16:14:55