如何通过Spring WebFlux实现Mono<List<Object>>的流式响应?
问题根因
当前实现的核心问题是collect(Collectors.toList())操作破坏了流式响应的基础:该操作会等待Flux<Result>中的所有元素全部返回后,才会将全量数据聚合为单个Mono<List<ResultDao>>一次性返回,哪怕配置了APPLICATION_STREAM_JSON_VALUE也不会触发分片输出,自然会因为接口处理时间过长触发AWS网关超时。
调整方案
方案1:标准流式JSON实现(推荐)
如果业务允许直接返回单个ResultDao的流式序列,直接修改返回类型为Flux<ResultDao>即可,这是WebFlux流式响应的原生支持用法:
@GetMapping(value= "/resultStream/{id}", produces = org.springframework.http.MediaType.APPLICATION_NDJSON_VALUE) public ResponseEntity<Flux<ResultDao>> getAllResultStream(@PathVariable String id) { Flux<Result> result = resultService.getAllResult(id); Flux<ResultDao> response = result.map(entity-> entity.convertToFlatObject()); return ResponseEntity.status(HttpStatus.OK).body(response); }
注:
APPLICATION_NDJSON_VALUE是标准换行分隔JSON流媒体类型,替代已废弃的APPLICATION_STREAM_JSON_VALUE,各类网关、客户端兼容性更好。
方案2:分片返回List结构的流式响应
如果业务要求每次返回的单元都是List<ResultDao>格式(比如按批次聚合数据),可以用buffer操作拆分Flux为多组List,返回Flux<List<ResultDao>>即可实现分片的List流式输出:
@GetMapping(value= "/resultListStream/{id}", produces = org.springframework.http.MediaType.APPLICATION_NDJSON_VALUE) public ResponseEntity<Flux<List<ResultDao>>> getAllResultListStream(@PathVariable String id) { Flux<Result> result = resultService.getAllResult(id); // 每100条数据聚合为一个List分片,可根据业务调整批次大小 Flux<List<ResultDao>> response = result.map(entity-> entity.convertToFlatObject()) .buffer(100); return ResponseEntity.status(HttpStatus.OK).body(response); }
注意事项
- 不要对返回的响应流做任何全量汇聚操作(
collect、reduce、block等),这类操作都会打断流式响应链路 - 如果使用AWS API网关做入口,需要关闭网关的响应缓存配置,避免网关缓存全量响应后才返回给客户端,导致流式失效
内容的提问来源于stack exchange,提问作者Sri Santhosh Natarajan
相关产品推荐
相关产品推荐

