WebFlux中如何正确提取Mono<List<T>>内容并向下游链路传递?
WebFlux 提取Mono内容的正确实现方案
响应式编程范式下不需要手动“提取”Mono中的值,只需要基于操作符构建处理流水线,最终由订阅者触发执行即可,Spring WebFlux框架原生支持响应式返回值,大部分场景下不需要手动调用subscribe。
下面分三种常见场景给出对应实现:
场景1:直接将List<Payload>作为接口返回值
直接修改Controller方法的返回值为Mono<List<Payload>>,将reader调用结果直接返回即可,Spring框架会自动订阅这个Mono,等数据就绪后自动序列化返回给客户端:
@PostMapping("/subset") public Mono<List<Payload>> read(@RequestBody RequestParams params){ return reader.read(params.getDate(), params.getAssetClasses(), params.getFirmAccounts(), params.getUserId(), params.getPassword()); }
场景2:需要先将List<Payload>传给下游服务处理再返回
根据下游服务的实现类型选择对应操作符:
- 下游服务是响应式实现(返回Mono/Flux类型):用
flatMap操作符 - 下游服务是同步实现:用
map操作符
示例代码:
@PostMapping("/subset") public Mono<ProcessResult> read(@RequestBody RequestParams params){ return reader.read(params.getDate(), params.getAssetClasses(), params.getFirmAccounts(), params.getUserId(), params.getPassword()) // 同步下游处理用map .map(payloadList -> { // 调用同步下游服务处理 ProcessResult result = syncDownstreamService.process(payloadList); return result; }) // 异步响应式下游处理用flatMap // .flatMap(payloadList -> reactiveDownstreamService.process(payloadList)) ; }
场景3:不需要等待处理完成,直接返回响应给客户端
这种场景是唯一需要手动调用subscribe的情况,要注意必须加上异常处理逻辑,避免异常丢失:
@PostMapping("/subset") public ResponseEntity<Void> read(@RequestBody RequestParams params){ reader.read(params.getDate(), params.getAssetClasses(), params.getFirmAccounts(), params.getUserId(), params.getPassword()) .subscribe( // 处理正常返回的结果 payloadList -> syncDownstreamService.process(payloadList), // 显式处理异常,避免异常被吞 error -> log.error("处理Payload数据失败", error) ); // 直接返回202接受状态,不等处理完成 return ResponseEntity.accepted().build(); }
注意事项
- 非必要不要调用
block()方法,如果必须兼容老的同步阻塞逻辑,需要先通过publishOn(Schedulers.boundedElastic())将执行切换到专门的阻塞线程池,再调用block,避免占用WebFlux的Netty核心Worker线程。 - 大部分场景下Spring WebFlux是流的最终订阅者,不要手动调用subscribe,否则会丢失请求上下文,异常也无法被框架的全局异常处理器捕获。
内容的提问来源于stack exchange,提问作者Simeon Leyzerzon
相关产品推荐
相关产品推荐

