Reactor特定异常重试实现:Webflux乐观锁场景下替代递归的方案
现有实现的问题
- 存在阻塞调用:
Thread.sleep是同步阻塞方法,在Webflux的非阻塞调度线程中使用会导致线程资源被占用,大幅降低服务吞吐量,不符合响应式编程的核心规范 - 递归风险:无边界的递归重试在冲突频率较高的场景下会出现栈溢出问题
- 异常链路断裂:嵌套调用
subscribe、方法返回void的写法,会导致操作结果、异常无法向上游传播,上层调用方无法感知操作状态,也无法做统一的异常处理 - 副作用操作不规范:将数据库保存操作放在
doOnNext副作用方法中,没有正确接入响应式流的信号链路,保存操作的异常可能无法被外层的异常捕获逻辑识别
最优实现方案
推荐直接使用Reactor内置的重试能力实现,完全符合响应式规范,无需自定义递归逻辑:
依赖引入
首先需要引入Reactor Extra组件提供的重试工具(Spring Boot项目通常已默认传递依赖,无需额外添加):
<dependency> <groupId>io.projectreactor.addons</groupId> <artifactId>reactor-extra</artifactId> </dependency>
代码实现
import reactor.util.retry.Retry; import java.time.Duration; public Mono<Void> updateStatus(String id, EventStatus status) { return eventRepository.findById(id) // 用flatMap将保存操作接入响应式流 .flatMap(eventDocument -> { eventDocument.setStatus(status); return eventRepository.save(eventDocument); }) // 配置重试规则 .retryWhen(Retry.backoff(3, Duration.ofSeconds(2)) // 仅对乐观锁异常重试 .filter(throwable -> throwable instanceof OptimisticLockingFailureException) // 可选:添加重试日志,方便排查问题 .doBeforeRetry(retrySignal -> log.info("乐观锁冲突,第{}次重试更新事件,id:{}", retrySignal.totalRetries() + 1, id) ) ) // 忽略返回的实体,返回空信号 .then(); }
方案优势
- 完全非阻塞:使用Reactor原生的异步延迟能力,没有任何阻塞逻辑,适配Webflux技术栈的性能要求
- 无递归风险:基于框架原生能力实现重试,不会出现栈溢出问题
- 重试规则可配置:支持自定义最大重试次数、延迟策略,示例中设置的是最多重试3次,每次间隔2秒的退避策略,相比固定延迟可以降低高并发下的冲突概率
- 异常链路完整:返回
Mono<Void>类型,上层调用方(比如消息监听器、其他服务调用方)可以正常感知操作成功/失败状态,方便后续做死信队列、告警等兜底逻辑 - 重试规则明确:仅对乐观锁异常进行重试,不会误重试其他业务异常
注意事项
- 最大重试次数需要根据业务场景合理设置,不推荐配置无限重试
- 超过最大重试次数后建议将异常向上抛出,在消息消费的入口处统一处理失败消息,比如投递到死信队列避免业务数据丢失
- 不要在业务方法内部调用
subscribe,订阅操作应该统一在最上层的入口(比如RocketMQ/Kafka消息监听器、Controller接口)执行
内容的提问来源于stack exchange,提问作者tsvlad
相关产品推荐
相关产品推荐

