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(/* 处理结果 */); } }
关键修改说明
- Mono实例缓存:使用
ConcurrentHashMap的computeIfAbsent原子方法,确保同一个ID只会创建一个查询流,并发请求共享该流,彻底避免重复API调用。 - 结果缓存:通过
cache()操作缓存流的最终结果,后续请求直接复用已有结果,无需重复走缓存查询流程。 - 线程池优化:对Redis/API这类可能阻塞的操作,使用
subscribeOn(Schedulers.boundedElastic())指定专用线程池,避免阻塞Reactor的IO线程。 - 内存泄漏防护:通过
doFinally在流完成后移除Mono缓存,避免大量ID导致内存占用过高(可根据业务场景选择是否保留)。
额外注意事项
- 确保
localCache是线程安全的实现,比如ConcurrentHashMap、GuavaLoadingCache或Caffeine Cache,避免并发写入/读取问题。 - 如果需要更灵活的缓存策略(比如自动过期、刷新),可以替换
cache()为CacheMono结合Caffeine等专业缓存框架实现。 - API调用后同步写入Redis,让后续请求可以从Redis获取数据,进一步降低API调用频率。
内容的提问来源于stack exchange,提问作者Techie97
相关产品推荐
相关产品推荐

