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

Reactive Streams:Spring WebFlux 订阅已有发布者问题咨询

解决Spring WebFlux中多请求复用同一数据发布者的问题

嘿,这个请求合并的场景我之前在项目里也遇到过,正好是Spring WebFlux里很常见的优化点——核心就是要避免对同一资源发起重复的并发请求,让后续请求直接“蹭”上正在执行的那个发布者就行。下面给你两个实用的实现方案,都是我实际用过的:

方案一:手动维护正在处理的请求(灵活可控)

这种方式适合你需要精细控制缓存逻辑的场景,核心思路是用一个并发Map来存正在执行的Mono实例,确保同一数据key下只会有一个获取请求在跑。

代码示例:

// 模拟你的缓存(可以换成Redis或其他分布式缓存)
private final Map<String, Data> localCache = new ConcurrentHashMap<>();
// 存正在处理中的请求Mono,避免重复调用
private final ConcurrentHashMap<String, Mono<Data>> inProgressFetches = new ConcurrentHashMap<>();

// 实际的耗时数据获取方法(比如调用外部API、查DB)
private Mono<Data> fetchActualData(String key) {
    return Mono.fromCallable(() -> {
        System.out.println("真的去拉数据啦: " + key);
        // 模拟耗时操作
        Thread.sleep(2000);
        return new Data(key, "新鲜数据");
    }).subscribeOn(Schedulers.boundedElastic());
}

public Mono<Data> getData(String key) {
    // 先查本地缓存,有就直接返回
    Data cached = localCache.get(key);
    if (cached != null) {
        return Mono.just(cached);
    }

    // 用computeIfAbsent保证同一key只会生成一个fetch请求
    return inProgressFetches.computeIfAbsent(key, k -> 
        fetchActualData(k)
            .doOnSuccess(data -> {
                // 成功后存缓存
                localCache.put(k, data);
                // 从正在处理的Map里删掉,避免内存占用
                inProgressFetches.remove(k);
            })
            .doOnError(err -> {
                // 失败也要删掉,不然后续请求会一直拿到失败的Mono
                inProgressFetches.remove(k);
            })
            // 缓存这个Mono,让后续订阅者共享结果
            .cache()
    );
}

这里的关键是ConcurrentHashMap的computeIfAbsent——它是原子操作,能确保同一时刻只有一个线程会为某个key创建fetchActualData的Mono,后面进来的请求直接复用这个已经在跑的Mono,完美避免重复请求。

方案二:用Spring Cache注解(懒人福音)

如果你项目里已经在用Spring Cache,那直接用@Cacheable就完事了,底层已经帮你封装好了请求合并的逻辑,根本不用自己写Map维护。

注意要确保你的缓存管理器支持Reactive,比如用Caffeine(需要引入spring-boot-starter-cache和caffeine依赖):

@Service
public class DataService {

    // 这里的cacheName和key可以根据你的需求调整
    @Cacheable(value = "data_cache", key = "#key")
    public Mono<Data> getData(String key) {
        return Mono.fromCallable(() -> {
            System.out.println("真的去拉数据啦: " + key);
            Thread.sleep(2000);
            return new Data(key, "新鲜数据");
        }).subscribeOn(Schedulers.boundedElastic());
    }
}

当多个请求同时打过来要同一个key的数据时,Spring Cache会自动让后面的请求等第一个请求完成,然后共享结果,完全不用你操心并发的问题,代码简洁到爆炸。

几个要注意的坑

  • 异常处理:不管用哪种方案,失败后一定要把正在处理的请求从Map里删掉(注解方式Spring会帮你处理),不然下次请求同一个key会直接拿到失败的Mono,永远拿不到数据。
  • 内存泄漏:手动维护Map的话,成功/失败后必须移除对应的key,不然Map会越来越大,占内存。
  • 缓存过期:如果需要数据过期,方案一可以用cache(Duration.ofMinutes(10))来设置缓存时长;方案二则可以在缓存配置里设置过期时间,比如Caffeine可以配置expireAfterWrite。

内容的提问来源于stack exchange,提问作者wild_nothing

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:38:54