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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:39:46