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

Webflux中'then'操作符未执行问题排查求助

Webflux响应式保存方法中then操作符导致后续逻辑不执行的原因排查

问题场景

编写Webflux响应式方法保存用户到数据库时,第一种实现使用then(reqMono)切换Mono后,后续的手机号重复校验及数据保存代码未执行;但将逻辑嵌套进第一个flatMap的第二种实现可正常运行。

问题代码

public Mono<Integer> saveAccount(Mono<AccountSaveRequest> requestMono) {
    
    var reqMono = requestMono.map(AccountSaveRequest::convertEntity);
    
    return reqMono
       .flatMap(req -> this.accountRepository.findByAccount(req.account())
                 .doOnNext(model -> {
                     if (req.id() == null || !model.id().equals(req.id())) {
                         throw new BusinessException("account exists");
                     }
                  })
        )
       .then(reqMono)
       .flatMap(req -> this.accountRepository.findByMobile(req.mobile())
                 .doOnNext(model -> {
                     if (req.id() == null || !model.id().equals(req.id())) {
                         throw new BusinessException("mobile exists");
                     }
                 })
        )
       .then(Mono.defer(() -> reqMono.flatMap(this.accountRepository::save)))
       .map(AccountEntity::id);
}

正常运行代码

public Mono<Integer> saveAccount(Mono<AccountSaveRequest> requestMono) {
    return requestMono.map(AccountSaveRequest::convertEntity)
       .flatMap(req -> this.accountRepository.findByAccount(req.account())
                  .doOnNext(model -> {
                      if (req.id() == null || !model.id().equals(req.id())) {
                          throw new BusinessException("account exists");
                      }
                  })
                  .switchIfEmpty(
                       Mono.defer(() -> this.accountRepository.findByMobile(req.mobile())
                           .doOnNext(model -> {
                                if (req.id() == null || !model.id().equals(req.id())) {
                                    throw new BusinessException("mobile exists");
                                }
                           })
                       )
                  )
                  .then(Mono.defer(() -> this.accountRepository.save(req)))
                  .map(AccountEntity::id)
        );
}

原因分析

1. then操作符+空Mono的连锁中断

当accountRepository.findByAccount(req.account())查询不到对应账号时,会返回空Mono(无元素发出,直接完成)。此时第一个flatMap的返回值就是这个空Mono,后续调用.then(reqMono)会切换到reqMono,但如果findByMobile也查询不到数据,同样返回空Mono。

Webflux中,如果Mono链全程只有完成信号、没有任何元素发出,最后.map(AccountEntity::id)因为没有元素可以映射,整个Mono会处于完成但无结果的状态,对外表现就是后续的保存逻辑没执行。

2. 冷订阅的额外干扰

reqMono是冷流,每次订阅都会重新执行requestMono.map()生成新的AccountEntity实例。第一种实现中reqMono被多次订阅,虽然这不是逻辑中断的核心原因,但会导致不必要的对象重复创建,也可能引发意料之外的状态问题。

3. 第二种实现的正确性逻辑

第二种实现将所有逻辑嵌套在同一个flatMap内部,用switchIfEmpty处理findByAccount返回空的情况,确保无论账号是否存在,都会执行手机号校验流程。整个链的信号是连续传递的,不会出现空Mono导致的无元素中断,最终能保证保存逻辑一定会被触发。

内容的提问来源于stack exchange,提问作者user8024151

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 23:43:08