如何让Observable重复执行10次并获取每次执行结果?
我之前也遇到过一模一样的需求——要重复调用日志API十次,不管每次成功还是失败都得拿到结果,一开始直接用repeat(10)踩了大坑:只要某次调用出错,整个Observable就直接终止了,后面的调用根本没机会执行。下面是我整理的可行方案,亲测有效:
核心思路
问题的关键在于:RxJava默认的repeat操作符只会在Observable成功完成时重复执行,一旦遇到错误事件就会直接终止流。所以我们得先把错误事件转换成普通的发射项,让流不会中断,再通过触发多次独立订阅来实现十次调用。
具体实现步骤
1. 包装API结果(可选但推荐)
先定义一个密封类统一包装成功和失败的结果,后续处理逻辑会更清晰:
sealed class ApiResult<out T> { data class Success<out T>(val data: T) : ApiResult<T>() data class Error(val exception: Throwable) : ApiResult<Nothing>() }
2. 改造原始API Observable
把Retrofit返回的Observable转换成能捕获错误并发射包装结果的流,这样错误不会终止整个流:
// 假设你的Retrofit接口返回Observable<LogResponse> fun getLogObservable(): Observable<ApiResult<LogResponse>> { return yourRetrofitApi.submitLog() .map { ApiResult.Success(it) as ApiResult<LogResponse> } .onErrorReturn { ApiResult.Error(it) } // 捕获错误,转换成Error类型结果发射 }
3. 触发10次调用并收集结果
用Observable.range(1, 10)生成10个触发信号,通过flatMap让每个信号都发起一次独立的API调用,这样不管前一次成功还是失败,下一次都会正常执行:
Observable.range(1, 10) .flatMap { callCount -> getLogObservable() .map { result -> Pair(callCount, result) } // 带上调用次数,方便区分每次结果 } .subscribeOn(Schedulers.io()) // Retrofit调用放在IO线程执行 .observeOn(AndroidSchedulers.mainThread()) // Android场景下切换回主线程处理结果 .subscribe( { (callCount, result) -> when (result) { is ApiResult.Success -> { println("第$callCount次调用成功:${result.data}") // 处理成功逻辑,比如更新UI或记录本地日志 } is ApiResult.Error -> { println("第$callCount次调用失败:${result.exception.message}") // 处理失败逻辑,比如记录错误信息 } } }, { // 这里不会走到,因为我们已经用onErrorReturn捕获了所有错误 } )
替代方案:用materialize操作符
如果你不想自定义结果类,可以用RxJava的materialize()操作符,它会把所有事件(成功的Next事件、失败的Error事件)都包装成Notification对象:
Observable.range(1, 10) .flatMap { callCount -> yourRetrofitApi.submitLog() .materialize() .map { notification -> Pair(callCount, notification) } } .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe { (callCount, notification) -> when { notification.isOnNext -> { val response = notification.value println("第$callCount次调用成功:$response") } notification.isOnError -> { val error = notification.error println("第$callCount次调用失败:${error?.message}") } } }
为什么直接用repeat(10)不行?
因为Retrofit的API Observable在调用失败时会发射Error事件,这个事件会直接终止整个Observable流,repeat操作符也就无法继续执行后续的重复逻辑。而我们通过onErrorReturn或materialize把错误转换成了正常的发射项,流会正常完成,就能保证10次调用都被执行。
内容的提问来源于stack exchange,提问作者Gurleen Sethi

