You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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();

关键调整说明

  1. 移动onErrorContinue位置:将onErrorContinue移至flatMap内部的每个save调用链中,仅针对当前元素的保存错误进行处理,不会干扰其他元素的正常流转,也不会吞噬重试成功后的信号。
  2. 移除全局retry():save方法内部已包含重试逻辑,全局retry()会导致整个Flux重新消费所有元素,违背“跳过失败元素”的需求。

内容的提问来源于stack exchange,提问作者Shakirov Ramil

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 22:41:10