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

Project Reactor中返回Mono<Void>的方法后如何使用flatMap衔接流程

Project Reactor 校验-存储-日志链路异常传递问题解决方案

业务场景说明

当前业务包含三类响应式组件:

  • 入参校验器:执行入参校验,无错误返回Mono.empty(),存在错误返回Mono.error(),返回值类型Mono<Void>
  • 数据库存储服务:负责数据持久化,返回值类型Mono<Item>
  • 操作日志组件:负责记录用户成功操作日志,返回值类型Mono<Void>

要求执行顺序:先做入参校验,校验通过后存数据库,最后记录操作日志,日志必须传入数据库保存完成后返回的Item实体。
预期执行规则:

  • 校验返回Mono.empty():继续执行后续调用链
  • 校验返回Mono.error():立即终止流程抛出异常,交由全局异常处理器处理

已尝试方案的问题

方案1:then()衔接存储

代码如下:

return inputValidator.validateFields(userId, projectId)
            .then(repository.save(item))
            .onErrorMap(RepoException.class, ex -> new UnexpectedError("Failed to save item", ex))
            .subscribeOn(Schedulers.boundedElastic())
            .doOnSuccess(n -> logService.logActivity(new Activity(adminId, n))
                    .subscribe());

该方案可在校验完成后触发存储,但会丢失校验阶段的错误。

方案2:flatMap()衔接存储

代码如下:

return inputValidator.validateFields(userId, projectId)
            .flatMap(v -> repository.save(item))
            .onErrorMap(RepoException.class, ex -> new UnexpectedError("Failed to save item", ex))
            .subscribeOn(Schedulers.boundedElastic())
            .doOnSuccess(n -> logService.logActivity(new Activity(adminId, n))
                    .subscribe());

该方案可正常传递校验异常,但因为校验通过返回Mono.empty(),flatMap不会触发后续存储逻辑。

核心问题说明

你对then()操作符的特性存在认知偏差:Project Reactor所有then系操作符都不会吞掉上游的error信号,只要上游发出error,then()会直接透传异常,绝不会执行传入的后续逻辑。你遇到的异常丢失问题,本质是validateFields方法实现有缺陷:校验失败时没有正确返回携带异常的Mono.error(),而是在方法内部吞掉异常最终返回了Mono.empty(),建议优先排查校验器代码。
另外你现有代码的日志写法存在严重问题:在doOnSuccess里单独调用subscribe()触发日志,会让日志逻辑脱离整个响应式链的管理,既无法保证线程上下文传递,也没法统一处理日志执行的异常,生产环境很容易出现日志丢失、上下文错乱的问题。

正确实现代码

优先修复校验器实现后,直接用以下写法即可满足所有需求:

return inputValidator.validateFields(userId, projectId)
        // 校验error直接透传,校验空完成后执行存储逻辑
        .then(repository.save(item))
        .onErrorMap(RepoException.class, ex -> new UnexpectedError("Failed to save item", ex))
        // 拿到存储返回的Item对象,衔接日志逻辑,保证整个链路在同一个响应式管道内
        .flatMap(savedItem -> logService.logActivity(new Activity(adminId, savedItem))
                .thenReturn(savedItem))
        .subscribeOn(Schedulers.boundedElastic());

如果你暂时无法修改校验器实现,可以用switchIfEmpty显式处理空信号,兼容异常丢失的场景:

return inputValidator.validateFields(userId, projectId)
        // 上游发出元素时执行存储
        .flatMap(v -> repository.save(item))
        // 上游空完成时同样执行存储,error信号直接透传
        .switchIfEmpty(repository.save(item))
        .onErrorMap(RepoException.class, ex -> new UnexpectedError("Failed to save item", ex))
        .flatMap(savedItem -> logService.logActivity(new Activity(adminId, savedItem))
                .thenReturn(savedItem))
        .subscribeOn(Schedulers.boundedElastic());

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 11:57:13