Observable未完成时判断isEmpty及缓存优先仓库实现求助
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
- Initial Cache Check:
getCacheSource().take(1)grabs the current cache state instantly (thanks to BehaviorSubject). We filter out empty data soinitialCacheCheckonly holds valid data or completes empty. - 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.
- 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 (orReplaySubjectwith buffer size 1). This guarantees new subscribers get the latest cache state immediately. - The
fetchFromNetwork()method should return aSingle<Data>orObservable<Data>that completes after emitting the fetched data. - When you call
cache.updateCache(), make suregetCacheSource()emits the new state—this is how ongoing updates reach the Presenter.
内容的提问来源于stack exchange,提问作者Jaeren Coathup

