RXJava-Retrofit链式API调用问题:PublishSubject的doOnNext未触发
核心原因分析
你的代码主要问题出在PublishSubject的热特性导致事件丢失,同时存在潜在的方法调用错误:
PublishSubject的事件丢弃特性
PublishSubject属于热Observable,仅会向订阅完成之后的观察者发送事件。当doOnSubscribe里执行this.getData(...).subscribe(subject)时,外部观察者还未完成对subject的订阅流程,若getData返回的Observable快速发射数据(比如Retrofit请求缓存命中、接口响应极快),这些onNext事件会被PublishSubject直接丢弃。如果getData的Observable紧接着触发onComplete,Subject会直接进入完成状态,后续外部观察者订阅后只会收到onComplete,完全收不到数据。方法调用参数不匹配
你在doOnNext里调用this.getData(taxonKey, regionId, null, nextOffset),但定义的getData方法仅接受3个参数(long id, String regionId, int offset),这里多传了一个null参数。如果存在对应的重载方法,可能该重载方法返回了不发射任何事件的Observable,直接导致流程无数据输出;若无重载方法,编译阶段就会报错。
修复方案
推荐放弃Subject实现,改用RxJava原生的concatMap+递归逻辑实现分页,这种方式更符合RxJava链式设计,彻底规避Subject的坑:
优化后的ObservationLoader代码
public Observable<PageObject<Observation>> getAllObservationDataByRegion(long taxonKey, String regionId) { // 从offset=0开始加载第一页 return loadPage(taxonKey, regionId, 0); } // 递归加载单页并推导后续请求 private Observable<PageObject<Observation>> loadPage(long taxonKey, String regionId, int offset) { return getData(taxonKey, regionId, offset) .concatMap(page -> { // 先发射当前页数据 Observable<PageObject<Observation>> currentPage = Observable.just(page); // 未到最后一页则递归加载下一页 if (!page.isEndOfRecords()) { int nextOffset = page.getOffset() + 1; Observable<PageObject<Observation>> nextPages = loadPage(taxonKey, regionId, nextOffset); return Observable.concat(currentPage, nextPages); } return currentPage; }); } private Observable<PageObject<Observation>> getData(long id, String regionId, int offset) { return this.api.getObservations(id, regionId, ObservationLoader.PAGE_LIMIT, offset) .subscribeOn(Schedulers.io()); } // HomeFragment订阅时统一指定线程 // observable.observeOn(AndroidSchedulers.mainThread()).subscribe(...);
方案优势
concatMap保证分页请求按顺序执行,前一页加载完成后才会发起下一页请求;- 完全基于RxJava链式操作,无需手动管理Subject的订阅和事件发射,从根源避免事件丢失;
- 递归逻辑清晰,当
endOfRecords=true时自动终止加载。
临时修复(保留Subject方式)
如果一定要沿用Subject实现,需改用ReplaySubject(会保存所有已发射事件,新订阅者能收到历史数据),同时修正方法调用参数:
public Observable<PageObject<Observation>> getAllObservationDataByRegion(long taxonKey, String regionId) { final ReplaySubject<PageObject<Observation>> subject = ReplaySubject.create(); return subject.doOnSubscribe(disposable -> { this.getData(taxonKey, regionId, 0).subscribe(subject); }) .doOnNext(observationPageObject -> { if (observationPageObject.isEndOfRecords()) { subject.onComplete(); } else { int nextOffset = observationPageObject.getOffset() + 1; // 修正参数,调用正确的getData方法 this.getData(taxonKey, regionId, nextOffset).subscribe(subject); } }) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()); }
内容的提问来源于stack exchange,提问作者Felix

