UniRx对Shared Observable使用Merge操作符时部分流无输出问题求解
问题根因
该问题由三个特性共同触发:
arrayInt.ToObservable()默认是同步发射所有值的冷Observable,一旦触发订阅会立刻把10个整数全部推送完成Share操作符会把冷流转为多播热流,只有发射事件时已完成订阅的观察者才能收到事件,错过的事件不会补发Observable.Merge会按传入参数的顺序依次订阅上游:先订阅observableNotEven时就触发了sharedObservable开始同步发射所有值,等所有值发射完毕才会开始订阅第二个参数observableFlatMap,此时所有事件已经推送完成,自然收不到任何值
解决方案
直接使用UniRx内置的ShareReplay操作符即可,无需修改其他业务逻辑,仅需要调整GetObservableInts方法中的一行代码:
把原代码中的.Share()替换为.ShareReplay()。ShareReplay默认会缓冲上游发射的所有历史值,在有新的订阅者接入时会先把所有缓冲值按顺序推送给新订阅者,再继续推送新的事件,同时它会在第一个订阅者到来时自动启动上游流,不需要手动调用Connect,既不会导致根Observable重复执行,也能完美适配你作为链路节点不需要单独控制生命周期的需求。
修改后的完整代码如下:
using System.Linq; using UniRx; using UniRx.Diagnostics; using UnityEngine; public class Share : MonoBehaviour { void Start() { PrintNumbers(); } private void PrintNumbers() { System.IObservable<int> sharedObservable = GetObservableInts(); var observableEven = sharedObservable. Where(x => x % 2 == 0) .Debug("Even"); var observableNotEven = sharedObservable. Where(x => x % 2 == 1) .Debug("NotEven"); var observableFlatMap = observableEven .Select(x => x * 10); _ = Observable.Merge(observableNotEven, observableFlatMap) .Subscribe(_number => Debug.Log(_number)) .AddTo(this); } private static System.IObservable<int> GetObservableInts() { var count = 10; var arrayInt = new int[count]; for (int id = 0; id != count; ++id) arrayInt[id] = id; var sharedObservable = arrayInt.ToObservable() .Debug("Array") .ShareReplay(); // 仅修改此处 return sharedObservable; } }
备选方案(不推荐)
如果你不需要缓冲历史值,也可以给上游添加异步调度器,让值的发射晚于两个流的订阅动作,比如将上游创建代码改为:
var sharedObservable = arrayInt.ToObservable(ThreadScheduler.Default) .Debug("Array") .Share();
该方案依赖调度器时序,稳定性不如ShareReplay。
内容的提问来源于stack exchange,提问作者HedgehogNSK
相关产品推荐
相关产品推荐

