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

Observable未完成时判断isEmpty及缓存优先仓库实现求助

Cache-First Repository Solution for Non-Completing Cache Observables

I get it—your cache Observable never completes because it needs to push ongoing updates to the Presenter, which breaks operators like isEmpty() that rely on the source finishing to emit a result. Let's fix this with an approach that checks the initial cache state upfront (without waiting for completion) while still listening for future cache changes.

Key Insight

If your cache uses a BehaviorSubject (or similar hot observable that emits the latest value to new subscribers), it will immediately send the current cache state when you subscribe. We can leverage this to check for valid data right away, then fall back to the network if needed.

Solution Code

Here's a RxJava implementation that meets your requirements:

public Observable<Data> getData() {
    // Step 1: Grab the initial cache state (emitted immediately if using BehaviorSubject)
    Maybe<Data> initialCacheCheck = getCacheSource()
            .take(1) // Get the first emission (current cache state)
            .filter(data -> !data.isEmpty()) // Only keep non-empty, valid data
            .firstElement(); // Convert to Maybe (emits data or completes empty)

    // Step 2: Fall back to network if cache is empty, then update cache
    Observable<Data> initialDataStream = initialCacheCheck
            .switchIfEmpty(
                fetchFromNetwork()
                    .doOnNext(fetchedData -> cache.updateCache(fetchedData)) // Update cache after successful network call
                    .toMaybe()
            )
            .toObservable();

    // Step 3: Combine initial data with ongoing cache updates (skip the first emission we already processed)
    return Observable.concat(
            initialDataStream,
            getCacheSource().skip(1)
    );
}

Breakdown of How This Works

  1. Initial Cache Check: getCacheSource().take(1) grabs the current cache state instantly (thanks to BehaviorSubject). We filter out empty data so initialCacheCheck only holds valid data or completes empty.
  2. Network Fallback: If the cache has no valid data, we switch to the network call. After fetching, we update the cache so subsequent cache emissions include the new data.
  3. Ongoing Updates: After the initial data (either from cache or network), we subscribe to the cache stream again (skipping the first emission we already handled) to push any future cache changes to the Presenter.

Important Notes

  • Ensure getCacheSource() returns a BehaviorSubject (or ReplaySubject with buffer size 1). This guarantees new subscribers get the latest cache state immediately.
  • The fetchFromNetwork() method should return a Single<Data> or Observable<Data> that completes after emitting the fetched data.
  • When you call cache.updateCache(), make sure getCacheSource() emits the new state—this is how ongoing updates reach the Presenter.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 11:20:13