Spring WebFlux中Flux赋值后doOnNext不执行的原因咨询
Reactor 中
doOnNext 未执行的本质原因分析 问题场景
在Spring Boot WebFlux项目中,发现doOnNext的执行情况不符合预期:只有Flux赋值给局部变量前的链式调用里的doOnNext会执行,赋值后单独调用的doOnNext完全无输出。
Controller代码如下:
@Controller public class ProblemController { @PostMapping("/upload") public Mono<ResponseEntity<String>> upload(ServerWebExchange exchange) { Flux<DataBuffer> body = exchange.getRequest().getBody() // 会打印输出! .doOnNext(db -> System.out.println("Bytes A: " + db.readableByteCount())); // 无任何输出... body.doOnNext(db -> System.out.println("Bytes B: " + db.readableByteCount())); return body .next() .map(db -> String.format("Uploaded %d bytes", db.readableByteCount())) .map(s -> new ResponseEntity<String>(s, HttpStatus.ACCEPTED)); } }
使用curl上传文件的命令:
curl --location --request POST 'http://localhost:8080/upload' \ --form 'file=@"/path/to/file.txt"'
执行后仅打印Bytes A,单独调用的Bytes B无输出;但如果把多个doOnNext都放在链式调用中,则全部能正常输出。
本质原因
这是Reactor响应式框架的两个核心特性导致的:
- 操作符的惰性执行:Reactor中的所有操作符都不会立即执行,只有当序列被订阅(subscribe)时,整个链上的操作才会被触发执行。
- 操作符返回新序列实例:每次调用操作符(比如
doOnNext),都会生成一个全新的Flux/Mono实例,而不是修改原有的实例。
对应到你的代码:
- 赋值给
body的是第一个Flux实例,它包含了Bytes A的doOnNext逻辑。后续返回链中使用了这个body,最终WebFlux会订阅这个实例,所以Bytes A的逻辑会被执行。 - 单独调用
body.doOnNext(...)时,生成了一个新的Flux实例,但这个实例既没有被赋值给任何变量,也没有被加入到最终的订阅链中——相当于这个新序列从未被订阅,因此它的doOnNext(Bytes B)永远不会触发。
如果要让Bytes B的逻辑执行,需要把这个新生成的Flux实例用在后续的订阅链中,比如:
Flux<DataBuffer> newBody = body.doOnNext(db -> System.out.println("Bytes B: " + db.readableByteCount())); return newBody .next() .map(db -> String.format("Uploaded %d bytes", db.readableByteCount())) .map(s -> new ResponseEntity<String>(s, HttpStatus.ACCEPTED));
内容的提问来源于stack exchange,提问作者Luke 10X
相关产品推荐
相关产品推荐

