You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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
根因分析
  1. 操作符使用违规:doOnError属于副作用回调操作符,仅用于日志打印、指标埋点等只读逻辑,禁止在回调内部抛出异常。在doOnError中直接抛出IllegalStateException会破坏Reactor内部的序列状态机,导致信号传递链路混乱。
  2. 操作符顺序不合理:onBackpressureBuffer属于流前置的背压控制操作符,放在多个生命周期回调之后,会在空流场景(无onNext信号直接下发onComplete)下产生信号重入。
  3. 版本兼容bug:2.7.12之前的Spring Boot版本对应的spring-data-mongodb响应式实现,存在空结果查询场景下,与Reactor背压操作符组合时重复下发完成信号的已知问题。
  4. 无doOnCancel日志属于正常表现:流全程没有收到取消信号,自然不会触发取消回调。
修复方案
  1. 规范副作用操作符的使用:移除doOnError内部的异常抛出逻辑,异常类型转换统一使用onErrorMap操作符实现,所有doOnXXX回调内仅做只读副作用操作,不抛出异常、不破坏流状态。
  2. 调整操作符顺序:将背压控制类操作符放在流链路最前端,紧接数据源返回的流实例,避免与下游生命周期回调产生信号干扰。
  3. 版本升级:如果使用低于2.7.12的Spring Boot版本,升级spring-data-mongodb到3.4.12及以上版本,修复空结果场景的信号重复下发bug。
  4. 防御性编程: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.01 22:21:37