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

如何让多次函数调用返回同一个正在执行的API请求Observable?

解决并发调用时重复发起API请求的问题

你的问题很典型——当多个订阅者同时调用你的getApiCall()方法时,因为缓存还未被填充,所有调用都会进入else分支,各自发起独立的API请求,导致重复请求。要解决这个问题,我们需要确保同一时间只有一个正在执行的API请求Observable被所有调用者共享,同时兼顾结果缓存。

核心思路

我们需要维护两个关键状态:

  1. 已缓存的API响应结果
  2. 当前正在进行的请求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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 13:32:32