Reactive Extensions Join非相交LeftDurationSelector失效及Switch疑问
问题解决:右Observable发射时携带左Observable最新值
最优方案:使用WithLatestFrom操作符
你要实现的场景——右序列每次发射值时,关联左序列的最新值——Rx已经提供了专门的操作符WithLatestFrom,这比用Join更直接且不易出错。直接替换你的代码即可:
var left1 = Observable.Interval(TimeSpan.FromSeconds(5)); var left2 = Observable.Interval(TimeSpan.FromSeconds(7)); var left = left1.Select(x => left2).Switch(); var right = Observable.Interval(TimeSpan.FromSeconds(1)); left.WithLatestFrom(right, (l, r) => (l, r)) .Distinct(x => x.r) .Subscribe(x => { Debug.WriteLine($"left: {x.l}, right: {x.r}"); });
你的Join代码出错原因
你最初的Join写法错误出在左持续选择器的逻辑:
_ => leftObservable.Publish().RefCount()
Join的左持续选择器的作用是:为左序列的每一个值,定义一个“有效周期Observable”——只有当这个周期Observable结束时,当前左值才会被移除,不再参与后续组合。你这里返回的是整个左序列的引用计数版本,而你的左序列是无限的(Interval生成),导致第一个左值永远不会失效,后续左值根本无法进入组合逻辑,自然没有输出。
如果非要用Join实现,正确的左持续选择器应该是让每个左值的有效周期持续到下一个左值发射,示例写法:
left.Join(right, _ => left.Skip(1).Take(1).IgnoreElements(), // 当前左值有效到下一个左值出现 _ => Observable.Empty<Unit>(), (l, r) => (l, r))
但显然WithLatestFrom是更优的选择。
关于Switch的语义澄清
你的left序列是left1.Select(x => left2).Switch(),Switch的核心语义是:每当源序列(left1)发射新的Observable(left2)时,立即取消订阅之前的Observable,转而订阅新的实例。所以left1每5秒发射一次,就会重新启动一个left2序列,每次新的left2都会从0开始计数——这是Switch的正常行为,你并没有误解它的作用,只是需要确认这种“定时重置左序列”的逻辑是否符合你的实际需求。
内容的提问来源于stack exchange,提问作者James B
相关产品推荐
相关产品推荐

