Spring WebFlux中Flux.doOnComplete()为何会被多次调用?
问题表现
- 基于Spring Boot WebFlux + 响应式MongoTemplate编写如下查询逻辑:
return mongoTemplateReactive.find(query, Document.class, COLLECTIONNAME). doOnNext(value -> ctr.incrementAndGet()). doOnCancel(() -> log.info("----Cancel Received----")). doOnError(ex -> { log.info("----Error----" + ex); throw new IllegalStateException(ex); }). onBackpressureBuffer(100000, BufferOverflowStrategy.DROP_OLDEST). doOnComplete(() -> { log.info("----OnComplete----"); });
- 运行时出现不符合响应式流规范的异常表现:
- 异常仅在MongoDB查询无匹配记录时触发,查询返回有效记录时所有逻辑运行正常
- 全程无
doOnCancel对应的日志输出,doOnComplete回调连续触发2次,日志输出如下:
----OnComplete---- ----OnComplete----
- 第一次
doOnComplete执行正常,第二次触发时回调内引用的外部对象已被置空,抛出NullPointerException
根因分析
- 操作符使用违规:
doOnError属于副作用回调操作符,仅用于日志打印、指标埋点等只读逻辑,禁止在回调内部抛出异常。在doOnError中直接抛出IllegalStateException会破坏Reactor内部的序列状态机,导致信号传递链路混乱。 - 操作符顺序不合理:
onBackpressureBuffer属于流前置的背压控制操作符,放在多个生命周期回调之后,会在空流场景(无onNext信号直接下发onComplete)下产生信号重入。 - 版本兼容bug:2.7.12之前的Spring Boot版本对应的spring-data-mongodb响应式实现,存在空结果查询场景下,与Reactor背压操作符组合时重复下发完成信号的已知问题。
- 无
doOnCancel日志属于正常表现:流全程没有收到取消信号,自然不会触发取消回调。
修复方案
- 规范副作用操作符的使用:移除
doOnError内部的异常抛出逻辑,异常类型转换统一使用onErrorMap操作符实现,所有doOnXXX回调内仅做只读副作用操作,不抛出异常、不破坏流状态。 - 调整操作符顺序:将背压控制类操作符放在流链路最前端,紧接数据源返回的流实例,避免与下游生命周期回调产生信号干扰。
- 版本升级:如果使用低于2.7.12的Spring Boot版本,升级spring-data-mongodb到3.4.12及以上版本,修复空结果场景的信号重复下发bug。
- 防御性编程:
doOnComplete回调内如果引用外部可变对象,增加非空判断,避免空指针异常。
修正后的参考代码:
return mongoTemplateReactive.find(query, Document.class, COLLECTIONNAME) // 背压操作前置 .onBackpressureBuffer(100000, BufferOverflowStrategy.DROP_OLDEST) .doOnNext(value -> ctr.incrementAndGet()) .doOnCancel(() -> log.info("----Cancel Received----")) // doOnError仅做日志打印,不抛异常 .doOnError(ex -> log.info("----Error----" + ex)) // 异常转换用专用操作符 .onErrorMap(IllegalStateException::new) .doOnComplete(() -> { log.info("----OnComplete----"); // 引用外部对象前先做非空校验 });
注意:所有doOnXXX系列操作符都是副作用操作,回调逻辑必须遵循只读、无异常抛出的原则,否则会直接破坏响应式流的状态一致性,出现信号重复、丢失、线程安全等不可预期的问题。
内容的提问来源于stack exchange,提问作者Pankaj Chouhan
相关产品推荐
相关产品推荐

