Reactor Mono主动缓存问题:定时刷新任务未实际生效求助
解决Reactor Mono主动缓存定时刷新的问题
这个问题我之前也碰到过,核心就是Mono的冷发布特性在搞鬼——它不会主动执行内部逻辑,必须有订阅者订阅时才会触发。所以你的定时刷新任务如果只是创建了Mono但没去订阅它,自然不会更新缓存。下面给你两种实用的实现方案:
方案一:定时任务主动订阅更新缓存
如果是在Spring环境下,用@Scheduled配合主动订阅是最直观的方式。我们可以维护一个线程安全的缓存变量,然后定时调用数据获取Mono并订阅,把新数据存入缓存:
代码示例
import reactor.core.publisher.Mono; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import java.util.concurrent.atomic.AtomicReference; @Component public class DataCacheService { // 用AtomicReference保证多线程环境下缓存的安全性 private final AtomicReference<Data> cachedData = new AtomicReference<>(); // 项目启动时初始化缓存,避免首次获取为空 @PostConstruct public void initCache() { updateCache().subscribe(); } // 每10分钟执行一次缓存更新(600000毫秒 = 10分钟) @Scheduled(fixedRate = 600000) public Mono<Void> updateCache() { return fetchDataFromSource() .doOnNext(data -> { cachedData.set(data); System.out.println("缓存已更新:" + data.getContent()); }) // 处理数据源异常,避免一次失败中断后续定时任务 .onErrorResume(e -> { System.err.println("数据获取失败,保留旧缓存:" + e.getMessage()); return Mono.empty(); }) .then(); // 转换为Mono<Void>标记任务完成 } // 对外提供获取缓存的方法 public Mono<Data> getCachedData() { return Mono.justOrEmpty(cachedData.get()); } // 模拟从数据源(API/数据库)获取数据的逻辑 private Mono<Data> fetchDataFromSource() { return Mono.fromCallable(() -> { // 替换为实际的数据源调用逻辑 return new Data("最新数据:" + System.currentTimeMillis()); }); } // 数据实体类示例 public static class Data { private String content; public Data(String content) { this.content = content; } public String getContent() { return content; } } }
关键说明
AtomicReference比volatile更稳妥,保证缓存操作的原子性和多线程可见性- 必须调用
subscribe()触发Mono执行,否则它的逻辑永远不会启动 onErrorResume用来兜底异常,避免一次数据源故障导致整个定时任务中断
方案二:用Reactor自身特性实现热流定时刷新
如果不想依赖Spring定时任务,可以用Reactor的repeatWhen和share()把冷Mono转为热流,让它主动定时刷新:
代码示例
import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; import java.time.Duration; import java.util.concurrent.atomic.AtomicReference; public class ReactiveDataCache { private final AtomicReference<Data> cachedData = new AtomicReference<>(); private Mono<Data> hotDataStream; public ReactiveDataCache() { initHotStream(); } private void initHotStream() { hotDataStream = fetchDataFromSource() .doOnNext(data -> { cachedData.set(data); System.out.println("缓存已更新:" + data.getContent()); }) // 每10分钟重复执行一次刷新 .repeatWhen(completed -> completed.delayElements(Duration.ofMinutes(10))) // 放在独立线程池执行,避免阻塞主线程 .subscribeOn(Schedulers.boundedElastic()) // 异常时保留旧缓存,1分钟后重试 .onErrorResume(e -> { System.err.println("数据获取失败,1分钟后重试:" + e.getMessage()); return Mono.just(cachedData.get()).delayElement(Duration.ofMinutes(1)); }) // 转为热发布流,多个订阅者共享同一个执行实例 .share(); // 主动订阅热流,启动定时刷新逻辑 hotDataStream.subscribe(); } public Mono<Data> getCachedData() { return Mono.justOrEmpty(cachedData.get()); } private Mono<Data> fetchDataFromSource() { return Mono.fromCallable(() -> new Data("最新数据:" + System.currentTimeMillis())); } public static class Data { private String content; public Data(String content) { this.content = content; } public String getContent() { return content; } } }
关键说明
share()将冷Mono转为热流,避免多个订阅者重复触发数据源调用repeatWhen会在每次流完成后,延迟指定时间重复执行,实现定时刷新- 初始化时必须主动调用
subscribe(),否则热流不会自动启动
核心总结
不管用哪种方案,主动订阅Mono是解决问题的关键——毕竟冷发布的特性决定了它不会自己“跑起来”。另外一定要注意缓存的线程安全和异常兜底,避免出现缓存失效或任务中断的情况。
内容的提问来源于stack exchange,提问作者Marcus Lindwall
相关产品推荐
相关产品推荐

