如何用Spring WebFlux实现嵌套循环并保证fetchResult执行顺序
Spring WebFlux 异步流改造:确保resultKeys生成后再调用fetchResult
在Spring WebFlux应用中实现了以下逻辑,但fetchResult(resultKeys)会在resultKeys生成完成前就被调用,而resultKeys是调用该方法的必要条件。需要将嵌套循环改造为符合Reactive编程模型的写法,保证fetchResult()仅在resultKeys完全生成后执行。
原同步代码
Map<MKey, Match> matchesByCom = intermediateObject.getMatchDetails(); Set<String> resultKeys = new HashSet<>(); intermediateObject.getNames().forEach(name -> { final Set<String> distinctIdsFinal = intermediateObject.getNameToIdMapping().get(name); matchesByCom.keySet().parallelStream().filter(mKey -> mKey.getName().equals(name)) .forEach(mKey -> { if (CollectionUtils.isEmpty(distinctIdsFinal) || (CollectionUtils.isEmpty(requestDTO.getLIds()))) { resultKeys.add(name + mKey.getCom()); } if (CollectionUtils.isNotEmpty(distinctIdsFinal)) { distinctIdsFinal.forEach(Id -> { resultKeys.add(name + mKey.getCom() + "_" + Id); }); } }); }); resultIdSetMono = fetchResult(resultKeys);
改造后的Reactive实现
要适配Spring WebFlux的异步非阻塞模型,需将同步遍历转换为Reactive流操作,利用Flux和Mono的操作符保证执行顺序:
Map<MKey, Match> matchesByCom = intermediateObject.getMatchDetails(); // 1. 将names集合转为异步流 Flux<String> namesFlux = Flux.fromIterable(intermediateObject.getNames()); // 2. 处理每个name,生成对应的resultKey流 Flux<String> resultKeysFlux = namesFlux.flatMap(name -> { Set<String> distinctIds = intermediateObject.getNameToIdMapping().get(name); // 筛选当前name对应的mKey,转为流处理 return Flux.fromIterable(matchesByCom.keySet()) .filter(mKey -> mKey.getName().equals(name)) .flatMap(mKey -> { Flux<String> currentKeys = Flux.empty(); // 添加基础key:name + com if (CollectionUtils.isEmpty(distinctIds) || CollectionUtils.isEmpty(requestDTO.getLIds())) { currentKeys = currentKeys.mergeWith(Flux.just(name + mKey.getCom())); } // 添加带Id的key:name + com + _ + Id if (CollectionUtils.isNotEmpty(distinctIds)) { currentKeys = currentKeys.mergeWith( Flux.fromIterable(distinctIds) .map(id -> name + mKey.getCom() + "_" + id) ); } return currentKeys; }); }); // 3. 收集所有resultKey到Set,确保生成完成后再调用fetchResult resultIdSetMono = resultKeysFlux .collect(Collectors.toSet()) .flatMap(this::fetchResult);
核心说明
- 所有集合遍历通过
Flux.fromIterable转为异步流,避免同步阻塞操作 - 使用
flatMap处理嵌套的流逻辑,保证每个元素异步处理完成后再推进流程 collect(Collectors.toSet())会等待所有resultKey生成完毕,将结果包装为Mono<Set<String>>- 最后通过
flatMap触发fetchResult,确保只有在resultKeys完全生成后才会执行该方法
内容的提问来源于stack exchange,提问作者Amulya M
相关产品推荐
相关产品推荐

