如何在Project Reactor的flatMap序列异常后继续执行后续操作
我是Project Reactor新手,遇到以下问题:现有一段使用flatMap的代码,链中第一个服务方法调用可能抛出异常,导致后续的服务方法调用无法执行,但我需要保证即使某一步出现异常,后续的服务方法仍能被调用。
实际场景为:从数据库获取记录,经serviceC处理后,调用serviceA的不同方法发送WEB、EMAIL等类型的Kafka异步消息(方法返回Mono<Void>)。目前我通过onErrorResume实现了临时解决方案,想询问该方案是否可行,或是否有更优的实现方式。
初始代码示例
// Message has a type of SomeClass flux.flatMap(message -> someService.someOperationReturningMono(...) .flatMap(i -> anotherService.operationWichCanProducerException(message.getProperty())) .flatMap(j -> anotherService.operationWichCanProducerException(message.getProperty()))) .then() .onErrorContinue(...) .subscribe(...);
临时解决方案代码示例
Mono<FakeDbRecord> mono = Mono.just(new FakeDbRecord("Propety")); mono.flatMap(message -> serviceC.someRandomLogic() .flatMap(i -> serviceA.methodA(message.getSomeProperty(), Types.WEB).then(Mono.just(message))) .onErrorResume(e -> { System.out.println("BOOOOOOOOOOM!!!!! " + e.getMessage()); return Mono.just(message); }) .flatMap(j -> serviceA.methodA(message.getSomeProperty(), Types.EMAIL).then(Mono.just(message))) .onErrorResume(e -> { System.out.println("BOOOOOOOOOOM!!!!! " + e.getMessage()); return Mono.just(message); })).subscribe();
方案可行性判断
你的临时方案完全可行。onErrorResume的设计初衷就是在序列抛出异常时,用一个正常的序列替代原序列继续执行,刚好匹配你“某步出错后不中断后续操作”的需求。你在每个可能抛出异常的服务调用后紧跟onErrorResume,确保前一步出错时仍能返回携带message的Mono,让后续的flatMap可以继续触发发送EMAIL消息的逻辑,完全符合业务要求。
更优实现方式
为了简化重复代码、提升可读性,可以针对发送消息的场景做以下优化:
1. 封装通用容错逻辑
把每个消息发送操作的容错处理抽成独立方法,避免重复编写onErrorResume:
// 封装单个消息发送的容错逻辑 private Mono<Void> sendMessage(FakeDbRecord message, Types type) { return serviceA.methodA(message.getSomeProperty(), type) .onErrorResume(e -> { System.out.printf("发送%s消息失败: %s%n", type, e.getMessage()); return Mono.empty(); // 用empty()表示异常已处理,不中断后续流程 }); } // 简化后的调用链 mono.flatMap(message -> serviceC.someRandomLogic() .then(sendMessage(message, Types.WEB)) // 先发送WEB消息,出错不中断 .then(sendMessage(message, Types.EMAIL))) // 再发送EMAIL消息,不受前一步影响 .subscribe();
这里用then替代flatMap,因为发送消息的方法返回Mono<Void>,我们只需要触发执行,不需要依赖前一步的返回值,then会忽略前序序列的结果,直接执行后续的Mono,逻辑更清晰。
2. 批量处理多个消息类型
如果需要发送的消息类型较多,可以用Mono.when结合流处理来批量执行,且每个子任务单独做容错:
import java.util.Arrays; import java.util.List; import java.util.stream.Collectors; // 定义要发送的消息类型列表 List<Types> messageTypes = Arrays.asList(Types.WEB, Types.EMAIL); mono.flatMap(message -> serviceC.someRandomLogic() .then(Mono.when( messageTypes.stream() .map(type -> sendMessage(message, type)) .collect(Collectors.toList()) ))) .subscribe();
这种方式适合并行发送多个消息的场景,单个消息发送的异常不会影响其他任务执行。
内容的提问来源于stack exchange,提问作者denstran

