如何统计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
相关产品推荐
相关产品推荐

