为何switchIfEmpty中的Mono未执行?Reactor代码问题排查
问题原因及解决方案
核心原因:类型不匹配导致流中断
你的代码存在类型不兼容的问题,直接导致functionThree中返回的Mono无法被正确订阅执行:
functionTwo的返回类型是Mono<List<Integer>>,但其中Mono.just(param).filter(...)返回的是Mono<Integer>,而switchIfEmpty要求传入的Mono类型必须与上游流的元素类型兼容(即Mono<? extends Integer>),但你传入的是Mono<List<Integer>>,类型完全不匹配。- 这种类型错误会导致编译失败,或运行时流抛出异常并终止,使得
reactiveStream.findById的Mono从未被真正订阅执行。 - 而
log.info("Valid")会输出,是因为调用functionThree(param)是同步执行的(构造Mono时立即执行),和流的订阅逻辑无关。
次要问题:不必要的同步调用
当前switchIfEmpty(functionThree(param))会立即调用functionThree,无论上游流是否为空,这会导致即使param不等于1时,Valid日志也会被输出,属于不必要的开销。
修复方案
- 修正类型匹配问题:在
functionTwo中,将Mono<Integer>转换为Mono<List<Integer>>,确保与switchIfEmpty传入的Mono类型一致:
Mono<List<Integer>> functionTwo(int param) { return Mono.just(param) .filter(p -> p != 1) // 将单个Integer转换为List<Integer>,匹配返回类型 .map(p -> Collections.singletonList(p)) .switchIfEmpty(functionThree(param)); }
- 优化懒加载调用:使用
Mono.defer包装functionThree的调用,确保只有上游流为空时才执行functionThree,避免不必要的同步操作:
Mono<List<Integer>> functionTwo(int param) { return Mono.just(param) .filter(p -> p != 1) .map(p -> Collections.singletonList(p)) // 延迟调用functionThree,仅在上游为空时执行 .switchIfEmpty(Mono.defer(() -> functionThree(param))); }
- 订阅时处理异常:在
SomeFunction的订阅逻辑中添加错误处理器,避免异常被静默丢弃,方便排查问题:
void SomeFunction(int param) { functionOne(param) .flatMap(p -> functionTwo(p)) .then(Mono.just(Boolean.TRUE.toString())) .subscribe( result -> log.info("Result: {}", result), error -> log.error("Stream error", error) // 添加错误日志 ); }
内容的提问来源于stack exchange,提问作者Spring boot progammer
相关产品推荐
相关产品推荐

