You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 10:45:28