如何在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
相关产品推荐
相关产品推荐

