Reactor带Retry调用链问题:onErrorContinue导致跳过重试后的值
问题:Reactor响应式调用链重试后值丢失与onErrorContinue冲突问题
在Reactor响应式调用链中执行带重试的repository save操作后,重试成功的值未被保留就进入后续流程。移除onErrorContinue后重试功能恢复正常,但会丢失该方法的异常处理能力——需要在save执行失败时跳过对应的SomeModel(适配Reactive Kafka监听器场景)。
操作流程
- 输入:Flux
- 在
flatMap中调用Model的repository save方法,并配置异常重试 - 核心需求:save执行失败时跳过对应的SomeModel,不影响其他元素处理
示例代码
主流程代码
Flux.fromIterable(s) // actually from kafka .flatMap(kk -> eventService.save(kk) .doOnNext(it -> record.receiverOffset().acknowledge())) .onErrorContinue((ex, record) -> { log.error("Exception during consumer", ex); if (record instanceof ReceiverRecord failedRecord && ex instanceof DeserializationException exception) { log.error("Failed record details {} and exception {}", record, exception); failedRecord.receiverOffset().acknowledge(); } else { log.error("Some unexpected error type! Cant move offset", ex); } }) .doOnError(e -> log.error("Exception during consumer", e)) .retry() .subscribe();
eventService.save方法(重试逻辑正常)
public Mono<SomeModel> save(SomeModel someModel) { return Mono.defer(() -> repository.save(someModel)) .retryWhen(Retry .backoff(1, Duration.ofMillis(100)) .maxBackoff(Duration.ofMillis(5000))); }
解决方案
问题根源
onErrorContinue的位置错误导致重试成功的信号被吞噬:原代码中onErrorContinue作用于整个Flux下游,当save方法内部抛出错误时,onErrorContinue会立即捕获并跳过当前元素,即使后续重试成功,信号也无法传递到doOnNext环节。
修改后的代码
Flux.fromIterable(s) // 实际来自Kafka .flatMap(kk -> eventService.save(kk) .doOnNext(it -> record.receiverOffset().acknowledge()) .onErrorContinue((ex, obj) -> { log.error("Exception during saving record", ex); if (obj instanceof ReceiverRecord failedRecord && ex instanceof DeserializationException exception) { log.error("Failed record details {} and exception {}", failedRecord, exception); failedRecord.receiverOffset().acknowledge(); } else { log.error("Unexpected error type! Cannot commit offset", ex); } }) ) .doOnError(e -> log.error("Global exception during consumer processing", e)) .subscribe();
关键调整说明
- 移动onErrorContinue位置:将
onErrorContinue移至flatMap内部的每个save调用链中,仅针对当前元素的保存错误进行处理,不会干扰其他元素的正常流转,也不会吞噬重试成功后的信号。 - 移除全局retry():
save方法内部已包含重试逻辑,全局retry()会导致整个Flux重新消费所有元素,违背“跳过失败元素”的需求。
内容的提问来源于stack exchange,提问作者Shakirov Ramil
相关产品推荐
相关产品推荐

