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

Spring WebFlux如何执行流方法并校验最终结果(无额外转换)

问题说明

现有返回值为Flux<String>的方法fluxFromFileStream,完成字符串处理、DTO转换后,需要调用transformAndSaveKpis(kpiHdfsDtoFlux)、transformAndSaveReports(kpiHdfsDtoFlux)两个方法将处理后的数据存入MongoDB。
当前核心实现代码如下:

public void transformAndSaveKpisAndReports(InputStream inputStream, DQSCJobName jobName) {
    fluxFromFileStream(inputStream)
            .flatMap(this::buildDTOFromLine)
            .map(kpiHdfsDto -> changeKpiTypeToFreezeIfJobNameIsFrozen(jobName, kpiHdfsDto))
            .cache()
            .transform(
                    kpiHdfsDtoFlux -> Flux.zip(transformAndSaveKpis(kpiHdfsDtoFlux), transformAndSaveReports(kpiHdfsDtoFlux))
            );
}

现存问题:

  • .transform()会返回无用的Flux对象
  • 需要校验整个流是否执行成功,执行过程中出现异常时直接在主方法transformAndSaveKpisAndReports中抛出
  • 此前通过判断流返回结果是否为null决定是否抛异常的实现不够规范优雅
    两个子方法实现代码如下:
private Flux<Kpi> transformAndSaveKpis(Flux<KpiHdfsDto> kpiHdfsDtoFlux) {
    return kpiHdfsDtoFlux
            .map(this::kpiHdfsDtoToKpiDocument)
            .collectList()
            .flatMapMany(kpis -> kpiRepository.insertAll(kpis));
}

private Flux<Report> transformAndSaveReports(Flux<KpiHdfsDto> kpiHdfsDtoFlux) {
    return kpiHdfsDtoFlux
            .flatMap(this::kpiHdfsDtoToReportDocument)
            .groupBy(Report::getType)
            .flatMap(reportList -> reportRepository.insertAll(reportList.collectList()));
}
实现方案

首先明确核心问题:Reactor流是懒加载的,当前代码仅完成了流组装,没有触发订阅/执行,两个存储逻辑实际上根本不会运行,这也是之前null判断逻辑无效的根本原因。
不需要使用transform()操作,直接合并两个存储流后触发执行即可,适配同步void返回场景的实现代码如下:

public void transformAndSaveKpisAndReports(InputStream inputStream, DQSCJobName jobName) {
    Flux<KpiHdfsDto> cachedDtoFlux = fluxFromFileStream(inputStream)
            .flatMap(this::buildDTOFromLine)
            .map(kpiHdfsDto -> changeKpiTypeToFreezeIfJobNameIsFrozen(jobName, kpiHdfsDto))
            .cache();

    Flux.zip(transformAndSaveKpis(cachedDtoFlux), transformAndSaveReports(cachedDtoFlux))
            .then()
            .block();
}

关键逻辑说明:

  • 将cache()后的流转为独立变量,保证上游文件解析、DTO转换逻辑仅执行一次,两个下游存储操作共享同一份处理后的数据,避免重复计算
  • 调用then()方法丢弃zip输出的无用Tuple2<Kpi, Report>元素,仅保留流的完成、异常状态信号,符合不需要返回业务数据的需求
  • 调用block()触发流实际执行,当前线程会阻塞直到两个存储操作全部完成,执行链路中任意位置抛出的异常都会在此处直接抛出,无需额外做null判断

注意:block()仅适用于响应式流与同步非响应式代码桥接的场景,当前主方法为void返回的同步方法,在此处调用block是合规用法,禁止在响应式调用链内部调用block。

如果运行环境为Spring WebFlux等不允许线程阻塞的场景,可将方法返回值改为Mono<Void>,移除block()直接返回流即可:

public Mono<Void> transformAndSaveKpisAndReports(InputStream inputStream, DQSCJobName jobName) {
    Flux<KpiHdfsDto> cachedDtoFlux = fluxFromFileStream(inputStream)
            .flatMap(this::buildDTOFromLine)
            .map(kpiHdfsDto -> changeKpiTypeToFreezeIfJobNameIsFrozen(jobName, kpiHdfsDto))
            .cache();

    return Flux.zip(transformAndSaveKpis(cachedDtoFlux), transformAndSaveReports(cachedDtoFlux))
            .then();
}

内容的提问来源于stack exchange,提问作者mr.Penguin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 11:18:16