Room的Reactive Streams问题:Flowable无法触发Reactive Streams Subscriber回调
这问题我之前也踩过坑,核心原因是Reactive Streams规范的背压机制在起作用,Room返回的Flowable严格遵循了这个规范,而你忽略了关键的一步操作。
为什么用org.reactivestreams.Subscriber没反应?
当你使用标准的org.reactivestreams.Subscriber订阅Flowable时,在onSubscribe()方法中拿到的Subscription对象,是用来控制上游数据发射的“开关”——你必须主动调用它的request()方法,告知上游你能接收多少数据,否则上游会一直处于等待状态,不会发射任何onNext/onError/onComplete事件。
看你的代码,onSubscribe()里只打印了日志,完全没处理Subscription的请求逻辑,这就是为什么只有onSubscribe()被执行,其他方法毫无动静的原因。
为什么用Consumer就能正常运行?
而当你使用RxJava的io.reactivex.functions.Consumer订阅时,RxJava的内部订阅逻辑已经帮你自动完成了request(Long.MAX_VALUE)的调用,相当于告诉上游“我能接收所有数据”,所以数据能正常发射到下游。
修复方案
只需要在onSubscribe()方法中添加request()调用即可,修改后的代码如下:
repository.getAgentInfo() .observeOn(AndroidSchedulers.mainThread()) .subscribe(new org.reactivestreams.Subscriber<AgentInfo>() { @Override public void onSubscribe(Subscription s) { Log.d("onSubscribe()", s.toString()); // 关键:请求上游发射所有数据,也可以指定具体数量比如request(1) s.request(Long.MAX_VALUE); } @Override public void onNext(AgentInfo agentInfo) { Log.d("onNext()", agentInfo.toString()); } @Override public void onError(Throwable t) { Log.d("onError()", t.getLocalizedMessage()); } @Override public void onComplete() { Log.d("onComplete()", ""); } });
额外说明
Room返回的Flowable是一个持续监听数据库变化的流,一旦调用request(Long.MAX_VALUE),它会在每次数据库中PartnerInfoEntity数据变化时,自动发射最新的数据到下游。如果你只需要获取一次数据,更推荐使用Single<AgentInfoInfoEntity>或者Maybe<AgentInfoInfoEntity>,这样不需要处理持续的数据流,也能避免不必要的资源消耗。
内容的提问来源于stack exchange,提问作者aleksey_bazhenov

