RxJava中如何在Completable完成前后差异化映射Observable
解决RxJava中Realm变更流与Retrofit请求状态的关联问题
嘿,这个场景我之前也踩过坑——merge确实不合适,因为它只是把两个流的事件混在一起,没法把网络请求的状态和Realm的变更结果对应起来。我给你推荐两种精准匹配需求的实现方式,核心思路都是跟踪Retrofit请求的状态,再将这个状态与Realm的持续变更流绑定。
方案一:用BehaviorSubject维护状态(推荐)
这种方式最直观,也最灵活,适合需要随时读取最新状态的场景:
1. 先定义状态枚举和数据包装类
首先我们需要语义化的状态标识,以及用来包装Realm数据和状态的类:
// 定义请求状态 enum RequestStatus { LOADING, // 请求进行中 SUCCESS, // 请求完成 ERROR // 请求失败 } // 包装Realm数据和状态的类 class RealmDataWrapper<T> { public final T realmResults; public final RequestStatus status; public RealmDataWrapper(T realmResults, RequestStatus status) { this.realmResults = realmResults; this.status = status; } }
2. 用BehaviorSubject跟踪请求状态
BehaviorSubject会保存最新的状态值,新订阅者一上来就能拿到当前的最新状态,完美适配Realm这种持续发射数据的流:
// 初始状态设为LOADING,代表请求还在进行中 BehaviorSubject<RequestStatus> statusSubject = BehaviorSubject.createDefault(RequestStatus.LOADING); // 订阅Retrofit的Completable,更新状态 yourRetrofitCompletable .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) // 按需调整线程 .subscribe( () -> statusSubject.onNext(RequestStatus.SUCCESS), // 请求成功,更新状态 throwable -> { statusSubject.onNext(RequestStatus.ERROR); // 这里可以加错误日志、Toast提示等逻辑 } );
3. 结合Realm流与状态流
用combineLatest操作符,每次Realm发射新的变更集,就用当前最新的请求状态来包装数据:
// 你的Realm变更Observable Observable<RealmResults<YourRealmModel>> realmChangeObservable = ...; // 合并两个流,生成带状态的包装数据Observable Observable<RealmDataWrapper<RealmResults<YourRealmModel>>> combinedObservable = Observable.combineLatest( realmChangeObservable, statusSubject, (realmData, currentStatus) -> new RealmDataWrapper<>(realmData, currentStatus) ); // 订阅最终的流 combinedObservable .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe( wrapper -> { // 这里根据状态处理逻辑: // LOADING时:显示本地数据+加载提示 // SUCCESS时:显示本地数据+隐藏加载提示/标记请求完成 // ERROR时:显示本地数据+错误提示 }, throwable -> { // 处理Realm流本身的错误 } );
方案二:纯流操作,不使用Subject
如果你不想用Subject,也可以直接通过流操作生成状态序列,同样能实现需求:
// 生成状态流:先发射LOADING,请求完成后发射SUCCESS,失败则发射ERROR,之后保持状态 Observable<RequestStatus> statusObservable = yourRetrofitCompletable .toObservable() .map(__ -> RequestStatus.SUCCESS) .onErrorReturn(__ -> RequestStatus.ERROR) .startWith(RequestStatus.LOADING) .concatWith(Observable.never()); // 让状态流不终止,持续保留最后状态 // 同样用combineLatest合并Realm流和状态流 Observable<RealmDataWrapper<...>> combined = Observable.combineLatest( realmChangeObservable, statusObservable, (data, status) -> new RealmDataWrapper<>(data, status) );
这个方案的好处是完全用流操作实现,没有外部状态,但灵活性稍差一些——如果后续需要手动更新状态,Subject方案会更方便。
为什么这两种方案比merge好?
- 它们都能精准关联Realm数据和请求状态:每个Realm的变更结果都会带上当前最新的请求状态,不会出现数据和状态不匹配的情况
- 语义更清晰:用
RequestStatus枚举代替Boolean,代码可读性更高 - 不需要对Completable做生硬的转换:不用把Completable转成发射Boolean的Observable,而是直接跟踪业务逻辑需要的状态
内容的提问来源于stack exchange,提问作者user3265561
相关产品推荐
相关产品推荐

