RxJava疑问:订阅共享Observable为何改变发射的条目
这其实是Rx系列库(比如RxJS)里冷Observable的典型特性导致的,我给你拆解清楚:
为什么注释掉stateChanges.subscribe()会丢失Request1?
默认情况下,Observable是冷流——也就是说,只有当有观察者主动订阅它时,它内部的生产者逻辑才会开始执行,而且每个新订阅都会从头完整执行一遍生产者逻辑。
如果你的Request1是绑定在stateChanges这个流的执行链路里(比如stateChanges的上游操作中生成了Request1,或者Request1的发射依赖stateChanges的触发逻辑),那当你没有订阅stateChanges时,这部分逻辑根本不会被触发。哪怕你的主流可能和stateChanges有关联,但如果主流的订阅没有覆盖到Request1所在的执行分支,就会出现Request1丢失的情况。
举个简化的类似场景,你就能明白:
const stateChanges = new Observable(subscriber => { console.log("Request1"); // 对应你的Request1输出 subscriber.next("state更新内容"); subscriber.complete(); }); // 主流合并了stateChanges和另一个固定内容 const mainStream = stateChanges.pipe(mergeWith(of("另一内容"))); mainStream.subscribe(console.log); stateChanges.subscribe(); // 注释掉这句,"Request1"就不会打印
这里的核心是:当你额外订阅stateChanges时,相当于触发了它的生产者逻辑一次,所以Request1被输出;而注释掉后,只有主流内部对stateChanges的订阅(如果有的话)才会触发,但如果你的结构中Request1的发射只和stateChanges的独立订阅绑定,那就会丢失。
如何实现无额外订阅时发射两个条目?
要解决这个问题,核心是把冷流转成热流,让多个观察者共享同一个流的执行实例,这样只要主流有订阅,流的生产者逻辑就会执行,所有关联的发射都会触发。常用的操作符有这些:
share():让多个观察者共享一个订阅,第一个观察者订阅时启动流,最后一个取消订阅时停止,是最常用的热流转换方式。publish()+connect():手动控制流的启动时机,调用connect()后不管有没有观察者,流都会开始执行,适合需要主动触发的场景。replay():会缓存流的发射值,后续订阅的观察者能收到之前的发射,适合需要获取历史值的场景。
还是用刚才的例子修改一下,就能实现无额外订阅输出两个条目:
// 关键:用share()把stateChanges转成热流 const stateChanges = new Observable(subscriber => { console.log("Request1"); subscriber.next("state更新内容"); subscriber.complete(); }).pipe(share()); const mainStream = stateChanges.pipe(mergeWith(of("另一内容"))); mainStream.subscribe(console.log); // 不需要额外订阅stateChanges,就能输出Request1、state更新内容、另一内容
这样处理后,主流的订阅会触发stateChanges的执行,Request1会被正常发射,同时主流合并的另一内容也会输出,完美实现你要的效果。
内容的提问来源于stack exchange,提问作者dimsuz

