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
相关产品推荐
相关产品推荐

