.NET ReactiveX异常现象:使用Subject<int>处理多IObservable<int>推送消息问题
问题根源:
Observable.Range参数理解错误 嘿,我一眼就揪出你代码里的问题了——你搞混了Observable.Range的第二个参数含义!
核心误解
Observable.Range的实际签名是这样的:
public static IObservable<int> Range(int start, int count)
第二个参数是要生成的元素总数,不是你以为的结束值!你代码里把end = i * 5 + 4作为第二个参数传入,这意味着每次生成的元素个数是end个,而不是从start到end的连续范围。
你的代码问题分析
拿循环里的几个例子拆解:
- 当
i=0时,start=0,end=4,调用Observable.Range(0,4)会生成4个元素:0,1,2,3(不是你预期的0-4) - 当
i=1时,start=5,end=9,调用Observable.Range(5,9)会生成9个元素:5,6,...,13(远超出你以为的5-9) - 当
i=8时,start=40,end=44,调用Observable.Range(40,44)会生成44个元素,从40一直到83——这就是你看到68、69这些意外数值的原因!
修正后的代码
把第二个参数改成固定的5(因为你每次想生成5个连续数):
using System; using System.Reactive.Linq; using System.Reactive.Subjects; using System.Linq; ISubject<int> valuesSubject = new Subject<int>(); valuesSubject.Subscribe(Console.WriteLine); foreach (int i in Enumerable.Range(0, 9)) { var start = i * 5; // 第二个参数改为5,代表生成5个连续元素 Observable.Range(start, 5).Subscribe(valuesSubject.OnNext); }
运行这段代码就能得到你预期的0到49的连续输出了。
顺便聊聊你最初的WebSocket客户端串联需求
其实你用Subject作为聚合序列的思路完全正确!因为Subject同时实现了IObservable和IObserver,每当你有新的WebSocket客户端消息序列时,只需要让这个序列订阅到你的聚合Subject上(就像示例里的Subscribe(valuesSubject.OnNext)),开发者只需要订阅一次这个Subject,就能自动接收所有后续新客户端的消息——完美解决了Concat需要重新订阅新序列的痛点!
内容的提问来源于stack exchange,提问作者Daniel Green
相关产品推荐
相关产品推荐

