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
相关产品推荐
相关产品推荐

