RxJava 2单Observable多Observer仅一个接收事件问题排查
首先咱们拆解下你遇到的问题:第二个订阅完全没输出,核心原因分为两个层面——代码里的语法/逻辑小错误,以及对RxJava冷热Observable的理解偏差,还有publish/replay操作符的误用。
一、先修正代码里的明显错误
这些错误可能直接导致第二个订阅无法正常执行:
getDataEmitter里的循环变量写错了:
你写的是emitter.onNext(token);,但循环里的变量应该是tokens[i],修正后:emitter.onNext(tokens[i]);(如果第一个订阅能正常运行,大概率是你贴代码时的笔误)
retrieve方法缺少返回类型,需要补充void:public void retrieve(String file) { // ... 原有代码 }第二个订阅的错误处理有语法问题:
原代码error -> System.out.println(error-2)会把error当成数字运算,应该改成打印字符串:error -> System.out.println("error-2")map(pair -> collect(request, pair))里的request变量未定义,这会导致编译失败,你需要确认这个变量的来源并补全。
二、核心问题:冷Observable的特性与publish/replay的正确用法
假设你已经修正了上面的语法错误,第二个订阅依然没输出的话,问题就出在你的sourceObservable是冷Observable:
RxJava中的冷Observable,每次有Observer订阅时,都会重新执行一遍整个上游的数据流逻辑(也就是getDataEmitter里的download、split循环等操作)。如果你的download或service.find(id)是一次性消耗资源的操作(比如下载后文件被读取耗尽、接口返回的数据流只能消费一次),那第二次订阅时上游就没有数据可发射,导致第二个Observer收不到任何事件。
你提到尝试了publish和replay但没解决,大概率是用法不对。正确的做法是把冷Observable转换成热Observable,让多个Observer共享同一份数据流:
方案1:用replay()实现数据重播(适合需要给晚订阅的Observer发送历史事件的场景)
public void retrieve(String file) { // 1. 用replay()将冷Observable转为可重播的ConnectableObservable // replay()会缓存所有发射的事件,给后续订阅的Observer重播 final ConnectableObservable<YourPairType> sourceObservable = getDataEmitter(file) .flatMap(id -> Observable.from(service.find(id)), Pair::of) .map(pair -> collect(request, pair)) .replay(); // 关键:转为可连接的Observable // 2. 先完成所有Observer的订阅 sourceObservable .flatMap(this::map) .map(this::fileFormat) .buffer(10) .subscribe(batched -> { System.out.println("b-1"); }, err -> { System.out.println("error-1"); }, () -> { System.out.println("completed-1"); }); sourceObservable .map(pair -> format(pair)) .subscribe(e -> { System.out.println("e-2" + e); }, error -> System.out.println("error-2"), () -> System.out.println("completed-2")); // 3. 手动触发数据流,此时两个Observer会共享同一份数据 sourceObservable.connect(); }
方案2:用publish()+autoConnect()实现自动触发(适合不需要缓存历史事件的场景)
public void retrieve(String file) { // 1. 用publish()+autoConnect(2),当订阅数达到2时自动启动数据流 final Observable<YourPairType> sourceObservable = getDataEmitter(file) .flatMap(id -> Observable.from(service.find(id)), Pair::of) .map(pair -> collect(request, pair)) .publish() .autoConnect(2); // 订阅数达标后自动触发,无需手动调用connect() // 2. 依次订阅两个Observer,当第二个订阅完成时,数据流自动启动 sourceObservable .flatMap(this::map) .map(this::fileFormat) .buffer(10) .subscribe(batched -> { System.out.println("b-1"); }, err -> { System.out.println("error-1"); }, () -> { System.out.println("completed-1"); }); sourceObservable .map(pair -> format(pair)) .subscribe(e -> { System.out.println("e-2" + e); }, error -> System.out.println("error-2"), () -> System.out.println("completed-2")); }
三、为什么你之前用publish/replay没生效?
大概率是你没有调用connect()(publish()/replay()返回的是ConnectableObservable,必须调用connect()才会触发上游数据流),或者在调用connect()之前只订阅了一个Observer,导致第二个Observer订阅时数据流已经执行完毕。另外,热Observable的核心是让多个Observer共享同一份上游数据,避免重复执行download这类耗时操作。
内容的提问来源于stack exchange,提问作者Bharath

