如何让多次函数调用返回同一个正在执行的API请求Observable?
解决并发调用时重复发起API请求的问题
你的问题很典型——当多个订阅者同时调用你的getApiCall()方法时,因为缓存还未被填充,所有调用都会进入else分支,各自发起独立的API请求,导致重复请求。要解决这个问题,我们需要确保同一时间只有一个正在执行的API请求Observable被所有调用者共享,同时兼顾结果缓存。
核心思路
我们需要维护两个关键状态:
- 已缓存的API响应结果
- 当前正在进行的请求Observable(用线程安全的方式管理,避免并发冲突)
当有调用请求时:
- 如果已有缓存,直接返回缓存结果
- 如果有正在进行的请求,返回这个请求的Observable,让所有调用者共享同一请求流
- 如果既无缓存也无正在进行的请求,创建新的请求Observable,同时将它标记为“正在进行”,请求完成后更新缓存并清除“正在进行”标记
实现方案一:用AtomicReference保证线程安全
这是最简洁且线程安全的实现方式,利用AtomicReference来原子性地管理当前请求Observable:
import java.util.concurrent.atomic.AtomicReference private var cache: ApiResponse? = null private val currentRequest = AtomicReference<Observable<ApiResponse>>() fun getApiCall(): Observable<ApiResponse> { // 优先返回缓存结果 cache?.let { return Observable.just(it) } // 尝试获取当前正在进行的请求 var ongoingRequest = currentRequest.get() if (ongoingRequest != null) { return ongoingRequest } // 创建新的请求Observable,处理响应并更新缓存 val newRequest = retrofitClient.doApiCall() .map { response -> doSomeStuffWithOutput(response) cache = response response } .doFinally { // 请求完成(成功/失败)后,清空当前请求标记 currentRequest.compareAndSet(newRequest, null) } .share() // 确保多个订阅者共享同一请求流 // 原子性地将新请求存入AtomicReference,避免并发冲突 return if (currentRequest.compareAndSet(null, newRequest)) { newRequest } else { // 若此时已有其他线程存入了请求,返回那个已存在的请求 currentRequest.get()!! } }
关键细节说明
AtomicReference确保了对当前请求Observable的操作是原子性的,避免多线程环境下的竞态条件.share()操作符让多个订阅者共享同一个请求流,不会重复执行API调用doFinally块在请求完成(无论成功或失败)后清空当前请求标记,这样如果请求失败,下次调用会重新发起请求
实现方案二:用BehaviorSubject管理请求流
如果你更倾向于用RxJava的Subject来实现,也可以用BehaviorSubject来缓存当前请求的流:
private var cache: ApiResponse? = null private var requestSubject: BehaviorSubject<ApiResponse>? = null fun getApiCall(): Observable<ApiResponse> { // 优先返回缓存结果 cache?.let { return Observable.just(it) } // 如果已有正在进行的请求Subject,返回它(用hide()防止外部干扰Subject状态) requestSubject?.let { return it.hide() } // 创建新的BehaviorSubject来承载请求流 val subject = BehaviorSubject.create<ApiResponse>() requestSubject = subject // 发起API请求,并将结果/错误发送到Subject retrofitClient.doApiCall() .map { response -> doSomeStuffWithOutput(response) cache = response response } .subscribe( { response -> subject.onNext(response) subject.onComplete() }, { error -> subject.onError(error) }, { // 请求完成后清空Subject引用 requestSubject = null } ) return subject.hide() }
关键细节说明
BehaviorSubject会将后续的订阅者自动连接到当前的请求流,确保所有订阅者收到同一结果hide()方法防止外部调用者直接调用Subject的onNext()/onError()等方法,保证流的安全性- 请求完成后清空
requestSubject,确保下次缓存为空时可以重新发起请求
为什么你的原始代码会出现重复请求?
你的原始代码中,当多个线程同时调用getApiCall()时,cache还未被设置,所有线程都会进入else分支,各自创建新的Observable.just(retrofitClient.doApiCall())——这意味着每个线程都会独立发起一次API请求,完全没有共享请求流的逻辑,自然会出现重复请求。
内容的提问来源于stack exchange,提问作者Michał Tajchert
相关产品推荐
相关产品推荐

