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

Java响应式编程异步执行问题:确保仅单次API调用

解决Project Reactor并发下重复API调用问题

问题根源

原代码中每个并发请求都会独立创建一个Mono流,即便使用subscribeOn(Schedulers.single()),也仅能让单个流的订阅操作在单线程执行,多个流仍会并行触发。当多个请求同时检测到本地缓存为空时,都会走到API调用环节,最终导致重复请求。

解决方案核心

通过缓存Mono实例而非仅缓存结果,让同一个ID的并发请求共享同一个流执行逻辑,确保只有第一个请求会触发完整的「本地缓存查询→Redis查询→API调用」流程,后续请求直接复用已有结果。

修改后的代码实现

import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

public class MetadataFetcher {
    // 线程安全的本地缓存,需确保实现线程安全(比如ConcurrentHashMap、Guava/Caffeine Cache)
    private final Map<String, Map<String, Object>> localCache;
    private final DataService cacheService; // Redis缓存服务
    private final DataService fallbackService; // API服务
    private final TransformationUtil transformationUtil;

    // 缓存每个ID对应的Mono实例,确保并发请求共享同一个流
    private final ConcurrentHashMap<String, Mono<Map<String, Object>>> monoCache = new ConcurrentHashMap<>();

    public MetadataFetcher(Map<String, Map<String, Object>> localCache, DataService cacheService, 
                           DataService fallbackService, TransformationUtil transformationUtil) {
        this.localCache = localCache;
        this.cacheService = cacheService;
        this.fallbackService = fallbackService;
        this.transformationUtil = transformationUtil;
    }

    public Mono<Map<String, Object>> enrichData(String id) {
        // 原子性获取或创建Mono实例,并发请求不会重复创建流
        return monoCache.computeIfAbsent(id, this::fetchMetadataWithFallback);
    }

    private Mono<Map<String, Object>> fetchMetadataWithFallback(String id) {
        return Mono.justOrEmpty(localCache.get(id))
                // 本地缓存为空,尝试从Redis获取
                .switchIfEmpty(Mono.defer(() -> fetchFromCache(id)))
                // Redis也为空,调用API获取
                .switchIfEmpty(Mono.defer(() -> fetchFromApi(id)))
                .defaultIfEmpty(transformationUtil.defaultEmptyMap())
                // 缓存结果,后续请求直接复用;可设置过期时间,比如cache(Duration.ofMinutes(5))
                .cache()
                // 流完成后清理Mono缓存,避免内存泄漏(可选,根据业务场景决定)
                .doFinally(signalType -> monoCache.remove(id));
    }

    private Mono<Map<String, Object>> fetchFromCache(String id) {
        return cacheService.getSymbolMetadata(id, Map.of())
                .flatMap(symbolData -> {
                    if (symbolData == null || symbolData.isEmpty()) {
                        return Mono.empty();
                    }
                    // 将Redis数据写入本地缓存
                    localCache.put(id, symbolData);
                    return Mono.just(symbolData);
                })
                // Redis操作如果是阻塞式,指定专用线程池避免阻塞IO线程
                .subscribeOn(Schedulers.boundedElastic());
    }

    private Mono<Map<String, Object>> fetchFromApi(String id) {
        return fallbackService.getSymbolMetadata(id, Map.of())
                .flatMap(symbolData -> {
                    if (symbolData == null || symbolData.isEmpty()) {
                        return Mono.empty();
                    }
                    // 将API数据写入本地缓存,同步写入Redis减少后续API调用
                    localCache.put(id, symbolData);
                    cacheService.putSymbolMetadata(id, symbolData); // 补充写入Redis,可选
                    return Mono.just(symbolData);
                })
                // API调用指定线程池
                .subscribeOn(Schedulers.boundedElastic());
    }

    public static void main(String[] args) {
        MetadataFetcher metadataFetcher = new MetadataFetcher(/* 注入依赖实例 */);

        // 调用示例:无需额外指定single线程池,内部已处理并发
        kafkaEventsFlux
                .flatMap(data -> metadataFetcher.enrichData(data.getId()))
                .subscribe(/* 处理结果 */);
    }
}

关键修改说明

  1. Mono实例缓存:使用ConcurrentHashMap的computeIfAbsent原子方法,确保同一个ID只会创建一个查询流,并发请求共享该流,彻底避免重复API调用。
  2. 结果缓存:通过cache()操作缓存流的最终结果,后续请求直接复用已有结果,无需重复走缓存查询流程。
  3. 线程池优化:对Redis/API这类可能阻塞的操作,使用subscribeOn(Schedulers.boundedElastic())指定专用线程池,避免阻塞Reactor的IO线程。
  4. 内存泄漏防护:通过doFinally在流完成后移除Mono缓存,避免大量ID导致内存占用过高(可根据业务场景选择是否保留)。

额外注意事项

  • 确保localCache是线程安全的实现,比如ConcurrentHashMap、Guava LoadingCache或Caffeine Cache,避免并发写入/读取问题。
  • 如果需要更灵活的缓存策略(比如自动过期、刷新),可以替换cache()为CacheMono结合Caffeine等专业缓存框架实现。
  • API调用后同步写入Redis,让后续请求可以从Redis获取数据,进一步降低API调用频率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 16:23:27