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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 21:53:15