资讯动态

深入解析 Reactive Extensions for .NET 的序列变换:Select、SelectMany、Cast 与 Materialize

发布时间:2026/9/29 7:40:22 来源:尧图企业网站定制
后端【免费下载链接】reactiveThe Reactive Extensions for .NET项目地址https://gitcode.com/gh_mirrors/re/reactive点击查看免费下载导读本文围绕 Rx.NETReactive Extensions for .NET官方入门教程《IntroToRx》第 6 章展开系统讲解序列变换Transformation的核心运算符Select、SelectMany、Cast以及Materialize/Dematerialize。我们将以 06_Transformation.md 为骨架结合 Rx.NET/Source/src/System.Reactive 的真实源码实现与 Tests.System.Reactive 测试用例解释每个运算符的语义、签名、典型用法、底层原理与易踩的坑。读完本文你将能熟练把IObservableTSource变换为IObservableTResult、正确区分拉取式与推送式SelectMany的行为差异、在需要时安全地做类型断言并用NotificationT以统一形式描述序列的完整生命周期。一、为什么需要变换从过滤到重塑我们消费的序列其值并不总是我们想要的格式有时源序列携带了多余信息需要只挑出关心的值过滤见 第 5 章 Filtering有时每个值需要被放大要么变成更丰富的对象要么展开成更多的值——这就是本章讨论的变换。变换运算符的共同特征是对源序列的每一个输入项产生一个或一组输出项并保持 Rx 的基本规则——输出的IObservableT同样遵循OnNext/OnError/OnCompleted契约。二、Select一对一的投影2.1 基本签名与语义最直接的变换运算符是Select。它接受一个函数将TSource类型的值转换为TResult类型的值从而把IObservableTSource变换为IObservableTResultIObservableTResult SelectTSource, TResult( this IObservableTSource source, FuncTSource, TResult selector)TSource与TResult不必不同。第一个例子对整数序列整体加 3类型不变IObservableint source Observable.Range(0, 5); source.Select(i i 3) .Dump(3)这里沿用了 第 5 章 Filtering 开头定义的Dump扩展方法输出如下3 -- 3 3 -- 4 3 -- 5 3 -- 6 3 -- 7 3 completed2.2 跨类型变换与匿名类型投影Select更常见的用途是改变元素类型。例如把整数转成字符Observable.Range(1, 5) .Select(i (char)(i 64)) .Dump(char);char -- A char -- B char -- C char -- D char -- E char completed也可以投影为匿名类型对象Observable.Range(1, 5) .Select(i new { Number i, Character (char)(i 64) }) .Dump(anon);anon -- { Number 1, Character A } anon -- { Number 2, Character B } anon -- { Number 3, Character C } anon -- { Number 4, Character D } anon -- { Number 5, Character E } anon completedSelect是 C# 查询表达式语法支持的标准 LINQ 运算符因此上面最后一个例子也可以写成var query from i in Observable.Range(1, 5) select new {Number i, Character (char) (i 64)}; query.Dump(anon);2.3 带索引的重载Rx 的Select还有一个重载selector函数接收两个参数第二个参数是元素在序列中的索引。当投影逻辑依赖元素位置时使用它。从源码看Rx 为这两个重载分别实现了Selector与SelectorIndexed两个内部 sink。索引重载在 Select.cs 中通过checked(_index)自增计数这意味着索引溢出会抛出异常并经由OnError送达下游// Select.cs 中的 SelectorIndexed public override void OnNext(TSource value) { TResult result; try { result _selector(value, checked(_index)); } catch (Exception exception) { ForwardOnError(exception); return; } ForwardOnNext(result); }值得注意Selector无索引版与SelectorIndexed都把_selector的调用包在try/catch中一旦投影函数抛出异常会被转换为OnError通知转发给观察者而不会破坏 Rx 的异常模型。这正是 Rx 运算符的通用范式用户回调中的异常一律进入错误通道而不是向上传播。三、SelectMany一对多的扁平化3.1 从嵌套序列到扁平序列Select一个输入只产出一个输出而SelectMany允许每个输入元素被变换成任意数量的输出。先看一个只用Select的例子Observable .Range(1, 5) .Select(i new string((char)(i64), i)) .Dump(strings);strings--A strings--BB strings--CCC strings--DDDD strings--EEEEE strings completed每个数字被转换为长度等于该数字的字符串。现在假设我们不把数字转成字符串而是转成IObservablechar——只需在构造完字符串后加.ToObservable()Observable .Range(1, 5) .Select(i new string((char)(i64), i).ToObservable()) .Dump(sequences);或者把投影表达式换成i Observable.Repeat((char)(i64), i)效果完全一致。输出并不太有用strings--System.Reactive.Linq.ObservableImpl.ToObservableRecursive1[System.Char] strings--System.Reactive.Linq.ObservableImpl.ToObservableRecursive1[System.Char] strings--System.Reactive.Linq.ObservableImpl.ToObservableRecursive1[System.Char] strings--System.Reactive.Linq.ObservableImpl.ToObservableRecursive1[System.Char] strings--System.Reactive.Linq.ObservableImpl.ToObservableRecursive1[System.Char] strings completed现在我们得到了一个可观察序列的可观察序列IObservableIObservablechar。把Select换成SelectMany再看Observable .Range(1, 5) .SelectMany(i new string((char)(i64), i).ToObservable()) .Dump(chars);结果变成单个IObservablecharchars--A chars--B chars--B chars--C chars--C chars--D chars--C chars--D chars--E chars--D chars--E chars--D chars--E chars--E chars--E chars completed输出顺序有些乱但仔细数一数每个字母出现的次数与之前发字符串时完全一致——只有一个AC出现三次E出现五次。SelectMany要求变换函数为每个输入返回一个IObservableT随后把所有这些结果合并回单个序列。3.2 与 LINQ to Objects 的对比拉取 vs 推送LINQ to Objects 版本的对应行为要规整得多Enumerable .Range(1, 5) .SelectMany(i new string((char)(i64), i)) .ToList()结果是[ A, B, B, C, C, C, D, D, D, D, E, E, E, E, E ]为什么会有这种差异这背后是IEnumerableT与IObservableT两种模型的根本区别IEnumerableT是拉取pull式的序列只在被要求时才产生元素。Enumerable.SelectMany按非常确定的顺序拉取——先问源IEnumerableint要第一个值把该值传给回调然后完整枚举回调返回的IEnumerablechar等它耗尽才去源里要第二个值如此循环。因此总是第一个内层序列的全部元素先出现然后是第二个内层序列的全部元素……。它之所以能这样推进一方面因为拉取模型让它可以自行决定处理顺序另一方面IEnumerableT的操作通常会阻塞直到拿到结果——上面的ToList不会在所有元素收集完毕前返回。Rx 不是这样的。第一消费者无法告诉源何时生产元素——源在准备好时自行发射。第二Rx 通常建模的是持续进行中的过程方法调用一般不会阻塞到完成。Rx 中的大多数操作会立即返回一个IObservableT或一个表示订阅的IDisposable值随后才产生。本文这个例子的 Rx 版本恰好属于少数特殊情况每个内层IObservablechar都在尽可能快地发射元素逻辑上所有内层序列并发进行。输出之所以显得交错是因为这些可观察源都试图尽快产出所有元素而它们交织的方式与 Rx 的调度器scheduler系统有关详见 第 11 章 Scheduling and Threading。调度器确保即使我们在建模逻辑上并发的过程Rx 的规则依然成立——SelectMany输出的观察者同一时刻只会收到一个元素。下图是产生上述交错输出的完整时序从源码也可以印证SelectMany的并发展开结构。在 SelectMany.cs 中外层 sink 为每个源元素调用_collectionSelector(value)创建内层序列并注册一个InnerObserver加入CompositeDisposable _group内层序列由独立的InnerObserver订阅各自向前转发OnNext。外层序列只有在_group.Count 0所有内层序列都完成时才向前转发OnCompleted——这正是SelectMany能并发展开多个内层序列、再把它们全部收拢回一个输出的实现基础。3.3 人为制造顺序Delay 与真正的严格顺序一个小的改动可以避免所有子序列同时竞争Observable .Range(1, 5) .SelectMany(i Observable.Repeat((char)(i64), i) .Delay(TimeSpan.FromMilliseconds(i * 100))) .Dump(chars);这里用Observable.Repeat替代了之前构造 string 再ToObservable的绕路写法——前者只是为了强调与 LINQ to Objects 例子的相似性实际 Rx 编程中不会那样写。现在输出与IEnumerableT版本一致chars--A chars--B chars--B chars--C chars--C chars--C chars--D chars--D chars--D chars--D chars--E chars--E chars--E chars--E chars--E chars completed下图展示了加上延迟后的时序——每个子序列成批地产出我们只是用死时间换来了分隔但请注意这些空隙只是为了讲解而引入的。如果真的追求严格按序处理实践中不应该这样用SelectMany——原因有二它并不能完全保证成功。如果缩短时间间隔最终会再次出现交错而 .NET 并非实时系统不存在任何可以保证顺序的安全时间值。如果确实需要第一个子序列的全部元素先于第二个子序列的任何元素出现有一个稳健的做法——用ConcatObservable .Range(1, 5) .Select(i Observable.Repeat((char)(i64), i)) .Concat()) .Dump(chars);当然这个例子不再使用SelectMany了Concat将在 第 9 章 Combining Sequences 中讨论。实践中的经验法则是当我们知道自己解包的是一个单值序列、或者不关心顺序、希望元素随到随取时才使用SelectMany。3.4 SelectMany 的重要性扇出再扇入如果你按顺序阅读本书其实已经在前面的章节见过两次SelectMany。第一个例子出现在 第 2 章 LINQ Operators and Composition 一节IObservableint onoffs from _ in src from delta in Observable.Return(1, scheduler) .Concat(Observable.Return(-1, scheduler) .Delay(minimumInactivityPeriod, scheduler)) select delta;注意查询表达式包含两个from子句时C# 编译器会将其编译为对SelectMany的调用。这个例子展示了 Rx 中一种常见模式——扇出fan out再扇入fan in为src产生的每个元素创建一个短暂的IObservableint即delta变量先产出1经过minimumActivityPeriod后产出-1从而统计最近发生的事件数。这是扇出部分——源序列的元素各自产生新的可观察序列而SelectMany的关键作用是把所有这些新序列扁平化回单一输出序列。第二个例子稍显不同来自 第 3 章 Representing Filesystem Events in Rx 一节它也把多个可观察源合并进一个可观察源但那份可观察序列的列表是固定的——FileSystemWatcher的每种事件各对应一个。该场景用的是Merge运算符第 9 章会讲到直接传入所有要合并的序列即可。然而由于代码还想做别的事延迟启动、自动释放、多订阅者共享同一源最终组合出的操作符链要求那段返回IObservableFileSystemEventArgs的合并代码作为变换步骤被调用。若用Select结果会是IObservableIObservableFileSystemEventArgs但代码结构决定了它只会产出单个IObservableFileSystemEventArgs双重包装的类型用起来非常不便。这正是SelectMany大显身手之处如果操作符组合引入了你不想要的多余序列套序列层SelectMany可以帮你解开一层。扇出再扇入与解开多余的可观察嵌套层这两种场景非常常见这让SelectMany成为 Rx 中的重要方法。在 Rx 的数学基础上SelectMany同样具有特殊地位——它是一种基础运算符许多其他 Rx 运算符都可以用它构建出来。附录 D 的 Recreating other operators withSelectMany一节 展示了如何用SelectMany实现Select和Where。在测试层面SelectManyTest.cs 覆盖了SelectMany_Complete、SelectMany_Complete_InnerNotComplete、SelectMany_Complete_OuterNotComplete、SelectMany_Error_Outer、SelectMany_Error_Inner、SelectMany_Dispose、SelectMany_Throw、SelectManyWithIndex_Index等大量场景分别验证外层/内层序列的完成、错误、释放与带索引重载的行为——如果你需要深入理解SelectMany的边界行为这些测试是很好的起点。四、Cast基于领域知识的类型断言C# 的类型系统并非全知全能。有时我们基于领域知识知道某个可观察源中值的更具体类型但类型的静态形式没有体现这一点。例如船舶广播的 AIS 消息如果消息类型是 3则其中包含导航信息。于是可以这样写IObservableIVesselNavigation type3 receiverHost.Messages.Where(v v.MessageType 3) .CastIVesselNavigation();Cast是标准 LINQ 运算符适用于我们确信集合中的元素是某个比类型系统能推断出的更具体的类型的场景。Cast与 第 5 章展示的OfType的关键区别在于对不合格元素的处理方式OfType是过滤运算符直接滤掉不是指定类型的元素Cast与普通 C# 强制转换表达式一样是断言我们声明源元素理应是该类型。若源产生一个与指定类型不兼容的元素Cast返回的序列将调用订阅者的OnError。用更基础的运算符重构两者差异一目了然// source.Castint(); 等价于 source.Select(i (int)i); // source.OfTypeint(); source.Where(i i is int).Select(i (int)i);源码证实了Cast的断言并上报错误语义。Cast.cs 的OnNext直接执行(TResult?)(object?)value强转任何失败都会捕获异常并ForwardOnError(exception)public override void OnNext(TSource value) { TResult? result; try { result (TResult?)(object?)value; } catch (Exception exception) { ForwardOnError(exception); return; } ForwardOnNext(result!); }CastTest.cs 中的Cast_NotValid、Cast_Error、Cast_Complete等用例即针对这些行为编写其中Cast_NotValid验证了类型不匹配时错误被正确转发。五、Materialize 与 Dematerialize把事件变成数据5.1 Materialize包装序列的完整生命周期Materialize运算符把IObservableT变换为IObservableNotificationT。源每产生一个元素它就提供一个NotificationT若源终止它还会产生一个最终的NotificationT说明序列是正常完成还是出错。这很有用因为它产出的对象完整描述了一条序列。如果想记录一个可观察源的输出以便日后重放通常直接用ReplaySubjectT它正是为此设计的但如果想重放之外还能做别的事——比如检查元素、甚至在重放前修改它们——就可能需要自己写代码存储元素。此时NotificationT的优势在于能用统一的方式表示源所做的一切无需单独记录序列是否终止、如何终止——这些信息就蕴含在最后一个NotificationT里。你甚至可以把它与ToArray结合用于单元测试得到一个NotificationT[]数组包含源所做一切行为的完整描述从而方便地断言序列的第三个元素是什么。Rx.NET 自己的源码就在大量测试中使用NotificationT。Materialize 一个序列可以看到被包装的值Observable.Range(1, 3) .Materialize() .Dump(Materialize);Materialize -- OnNext(1) Materialize -- OnNext(2) Materialize -- OnNext(3) Materialize -- OnCompleted() Materialize completed注意源序列完成时materialized 序列会先产出一个OnCompleted通知值然后自身才完成。NotificationT是抽象类有三个实现OnNextNotificationOnErrorNotificationOnCompletedNotificationNotificationT暴露四个公共属性用于检查Kind、HasValue、Value、Exception。只有OnNextNotification的HasValue为true且Value有实际意义同样只有OnErrorNotification的Exception有值。Kind属性返回一个枚举告诉你该使用哪些方法public enum NotificationKind { OnNext, OnError, OnCompleted, }这些成员在 Notification.cs 中有完整定义NotificationKind枚举、Kind、HasValue、Value、Exception等工厂方法Notification.CreateOnNextT、Notification.CreateOnErrorT、Notification.CreateOnCompletedT也在同一文件Notification.cs中提供。下面这个例子产生一条出错的序列。注意 materialized 序列的最后一个值是OnErrorNotification并且 materialized 序列本身不会出错而是正常完成var source new Subjectint(); source.Materialize() .Dump(Materialize); source.OnNext(1); source.OnNext(2); source.OnNext(3); source.OnError(new Exception(Fail?));Materialize -- OnNext(1) Materialize -- OnNext(2) Materialize -- OnNext(3) Materialize -- OnError(System.Exception) Materialize completed源码清楚地展示了这一语义。Materialize.cs 中OnNext(value)→ 转发Notification.CreateOnNext(value)OnError(error)→ 转发Notification.CreateOnErrorTSource(error)然后ForwardOnCompleted()错误被物化为数据序列本身正常完成OnCompleted()→ 转发Notification.CreateOnCompletedTSource()然后ForwardOnCompleted()。5.2 Dematerialize解包通知对序列做分析或日志记录时Materialize 非常方便。要解开一条 materialized 序列使用Dematerialize扩展方法——它只对IObservableNotificationTSource起作用。Dematerialize.cs 的实现按NotificationKind分发把三种通知分别还原为对应的OnNext、OnError、OnCompletedpublic override void OnNext(NotificationTSource value) { switch (value.Kind) { case NotificationKind.OnNext: ForwardOnNext(value.Value); break; case NotificationKind.OnError: ForwardOnError(value.Exception!); break; case NotificationKind.OnCompleted: ForwardOnCompleted(); break; } }DematerializeTest.cs 提供了Materialize_Dematerialize_Never/Empty/Return/Throw等成对测试验证物化—解包往返后的序列行为与源一致MaterializeTest.cs 则覆盖了Materialize_Never/Empty/Return/Throw各场景。六、小结本章的变换运算符有一个共同特征为源序列的每个输入项产生一个SelectMany则是产生一组输出。Select提供最直接的一对一投影并带索引重载SelectMany让每个输入展开为任意数量的输出并扁平化回收是扇出/扇入模式与解开嵌套序列的关键工具同时要留意它本质上的并发交织行为与拉取式 LINQ 的差异Cast以断言 失败走OnError的方式做类型升级区别于过滤式的OfTypeMaterialize/Dematerialize把序列的所有事件包括终止方式统一表示为NotificationT数据极大方便了记录、分析与测试。接下来我们将在后续章节研究能够结合源序列中多个元素信息的运算符。赞分享后端【免费下载链接】reactiveThe Reactive Extensions for .NET项目地址https://gitcode.com/gh_mirrors/re/reactive点击查看免费下载相关推荐Reactive Extensions for .NET 项目推荐Reactive Extensions for .NET 项目推荐 项目基础介绍和主要编程语言 Reactive Extensions for .NET简称后端Reactive Extensions for .NET 项目教程Reactive Extensions for .NET 项目教程 1. 项目的目录结构及介绍 Reactive Extensions for .NET简称后端Kirki配置系统详解打造个性化WordPress主题设置界面Kirki配置系统详解打造个性化WordPress主题设置界面 Kirki是一款强大的WordPress配置框架它能帮助开发者轻松创建直观且功能丰富的主题设开发工具UI库/组件上一篇如何使用SSTImap从入门到精通的交互式SSTI漏洞测试教程下一篇Vue FlatPickr 组件常见问题解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑