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

如何统计RxJava 2已执行的重试次数?是否有doOnRetry获取重试次数?

统计RxJava 2的重试总次数及替代doOnRetry的方案

嘿,这个问题问得很接地气!在RxJava 2里确实没有直接提供doOnRetry这个方法,但咱们有好几种灵活的方式来统计重试次数,下面给你一步步讲清楚:

一、核心方案:用retryWhen统计重试次数

这是最常用也最灵活的方式,retryWhen允许我们在每次重试触发前自定义逻辑,刚好可以用来计数。

由于RxJava的异步特性,计数变量必须是线程安全的,所以推荐用AtomicInteger来追踪次数,示例代码如下:

import io.reactivex.Observable;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public class RetryCountDemo {
    public static void main(String[] args) throws InterruptedException {
        AtomicInteger retryCount = new AtomicInteger(0);

        Observable.create(emitter -> {
            // 模拟可能抛出异常的业务逻辑
            System.out.println("执行核心任务...");
            emitter.onError(new RuntimeException("模拟业务异常"));
        })
        .retryWhen(errors -> errors.flatMap(error -> {
            int currentRetry = retryCount.incrementAndGet();
            System.out.println("触发第 " + currentRetry + " 次重试");
            
            // 可添加重试终止条件,比如最多重试3次
            if (currentRetry < 3) {
                // 重试前延迟1秒,避免频繁请求
                return Observable.timer(1, TimeUnit.SECONDS);
            }
            // 超过重试次数,传递错误终止流
            return Observable.error(error);
        }))
        .subscribe(
            success -> System.out.println("任务执行成功!"),
            fail -> System.out.println("任务最终失败,累计重试次数:" + retryCount.get())
        );

        // 阻塞主线程等待结果(仅示例用,实际业务按需处理)
        Thread.sleep(5000);
    }
}

二、模拟doOnRetry的效果:自定义扩展操作符

如果你的项目里经常需要监听重试事件,可以封装一个类似doOnRetry的扩展方法,方便复用:

import io.reactivex.Observable;
import io.reactivex.ObservableTransformer;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;

public class RxRetryExtensions {
    // 自定义doOnRetry操作符,接收一个消费重试次数的回调
    public static <T> ObservableTransformer<T, T> doOnRetry(Consumer<Integer> retryCallback) {
        AtomicInteger retryCounter = new AtomicInteger(0);
        return upstream -> upstream.retryWhen(errors -> errors.flatMap(error -> {
            int currentRetry = retryCounter.incrementAndGet();
            try {
                // 触发重试回调,传递当前重试次数
                retryCallback.accept(currentRetry);
            } catch (Exception e) {
                // 如果回调抛出异常,终止流并传递该异常
                return Observable.error(e);
            }
            // 允许继续重试(也可在此添加重试条件)
            return Observable.just(currentRetry);
        }));
    }
}

使用这个扩展方法的示例:

Observable.create(emitter -> {
    System.out.println("执行核心任务...");
    emitter.onError(new RuntimeException("模拟异常"));
})
.compose(RxRetryExtensions.doOnRetry(retryNum -> 
    System.out.println("监听到第 " + retryNum + " 次重试")
))
.retry(3) // 限制最多重试3次
.subscribe(
    success -> System.out.println("成功完成"),
    fail -> System.out.println("最终失败")
);

三、注意事项

  • 必须用线程安全的计数类(比如AtomicInteger),因为RxJava的流可能在不同线程执行,普通int会有线程安全问题;
  • 如果只是简单限制重试次数并统计,retry(long count, Predicate<Throwable>)也可以配合计数,但灵活性不如retryWhen;
  • RxJava 2确实没有内置的doOnRetry方法,但通过上面的方式完全可以实现相同甚至更灵活的效果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:58:32