Project Reactor如何优雅实现带条件的Mono交叉中断调用
你完全不需要手动维护AtomicReference存储Disposable来控制取消逻辑,用Project Reactor内置的操作符组合就能优雅实现所有需求,代码量远低于手写实现,也不会出现InterruptedException相关的错误日志。
实现思路
核心利用Mono.firstWithValue的原生特性:
- 并行订阅所有传入的上游Mono
- 只要收到任意一个上游发出的非空有效值,立刻向下游传递结果,同时取消其他所有还在运行的上游
- 如果某个上游提前完成但没有发出有效值(返回空Mono),会继续等待其余上游,不会提前终止流程
你只需要对每个餐厅的返回结果做过滤:仅保留返回值为"Beer!"的有效结果,不提供啤酒的餐厅返回会被转为空完成信号,刚好匹配firstWithValue的等待逻辑;最后加一个空值兜底,就能覆盖两家都不提供啤酒的场景。
可直接运行的实现代码
import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; public class ReactorBeerDemo { public static void main(String[] args) { // 处理Sos Diner的调用:只保留有啤酒的结果,绑定弹性调度器实现并行 Mono<String> sosDinerMono = contactSosDiner() .filter(response -> "Beer!".equals(response)) .doOnCancel(() -> System.out.println("Sos Diner 调用已取消")) .subscribeOn(Schedulers.boundedElastic()); // 处理Burger King的调用,逻辑同上 Mono<String> burgerKingMono = contactBurgerKing() .filter(response -> "Beer!".equals(response)) .doOnCancel(() -> System.out.println("Burger King 调用已取消")) .subscribeOn(Schedulers.boundedElastic()); // 核心逻辑:取第一个返回有效啤酒结果的源,空结果兜底 Mono.firstWithValue(sosDinerMono, burgerKingMono) .map(source -> "成功在" + source + "买到啤酒") .defaultIfEmpty("两家餐厅都不提供啤酒,没买到") .doOnNext(System.out::println) .block(); // 仅做演示使用,正式业务代码中避免随意调用block } public static Mono<String> contactSosDiner() { return Mono.fromCallable(() -> { Thread.sleep(1000); return "Beer!"; }); } public static Mono<String> contactBurgerKing() { return Mono.fromCallable(() -> { Thread.sleep(1500); return "No Beer, only whopper"; }); } }
场景验证
调整两个方法的sleep时长和返回值,可以覆盖所有你提到的分支:
- Sos Diner先返回且有啤酒:立刻取消Burger King的调用,输出买到啤酒的结果
- Sos Diner先返回但无啤酒:不触发取消,继续等待Burger King的返回
- Burger King先返回且有啤酒:立刻取消Sos Diner的调用
- Burger King先返回但无啤酒:继续等待Sos Diner的结果
- 两家都无啤酒:输出兜底的没买到啤酒的提示
运行默认参数的代码,输出如下,没有任何异常日志:
Burger King 调用已取消 成功在Sos Diner买到啤酒
实现优势
你之前手写的逻辑本质是重复实现了Reactor内置的上游调度、取消传播逻辑,不仅代码冗余,还没有正确处理阻塞任务被取消时抛出的InterruptedException,才会触发onErrorDropped的全局错误日志。用内置操作符时,Reactor框架会自动处理取消流程中的异常捕获、资源释放逻辑,不需要手动处理这些底层细节。如果后续需要新增更多餐厅作为候选,只需要把新的餐厅Mono按同样规则加过滤、绑定调度器后传入Mono.firstWithValue即可,扩展性更强。
内容的提问来源于stack exchange,提问作者tjamarant
相关产品推荐
相关产品推荐

