ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Rx.NET 数据过滤操作符完全指南:从 Where 到 DistinctUntilChanged 的事件裁剪艺术

Rx.NET 数据过滤操作符完全指南:从 Where 到 DistinctUntilChanged 的事件裁剪艺术 后端【免费下载链接】reactiveThe Reactive Extensions for .NET项目地址https://gitcode.com/gh_mirrors/re/reactive点击查看免费下载本文基于 Rx.NET/Documentation/IntroToRx/05_Filtering.md 展开围绕 .NET 响应式扩展Reactive Extensions for .NET即仓库中的Rx.NET项目的过滤Filtering操作符族系统讲解如何从高频、海量的事件流中裁剪出真正有价值的部分。读完本文你将掌握Where、IgnoreElements、OfType、FirstAsync/LastAsync/SingleAsync系列、Take/Skip系列、TakeWhile/SkipWhile、TakeUntil/SkipUntil、Distinct/DistinctUntilChanged的语义、签名、边界行为与源码实现位置并能在自己的 Rx.NET 应用中直接套用文中示例完成实战过滤。过滤的意义从数据洪流到高价值事件Rx 的核心价值在于把海量、持续到达的事件流转化为更高层次的洞察。这一过程通常伴随着数据量的缩减当低流量事件流中的每个个体事件平均携带更多信息时少量事件反而比大量事件更有用。Rx 提供了多种操作符来实现这种降噪其最简单的机制就是过滤掉我们不想要的事件。为了便于展示后续示例原文档定义了一个Dump扩展方法它订阅任意IObservableT将源产生的每条通知OnNext的值、OnError的异常、OnCompleted打印到控制台并接受一个name参数用于标记事件来源从而在订阅多个源的示例中区分输出。该辅助方法位于Rx.NET/Documentation/IntroToRx/05_Filtering.mdpublic static class SampleExtensions { public static void DumpT(this IObservableT source, string name) { source.Subscribe( value Console.WriteLine(${name}--{value}), ex Console.WriteLine(${name} failed--{ex.Message}), () Console.WriteLine(${name} completed)); } }按内容过滤Where对序列施加过滤是最常见的需求而 LINQ 中最直接的过滤器就是Where。与 LINQ 惯例一致Rx 也以扩展方法形式提供该操作符其签名如下IObservableT WhereT(this IObservableT source, FuncT, bool predicate)注意source参数与返回值的元素类型同为T因为Where不修改元素本身它只会剔除部分元素保留的元素原样透传。以下示例用Where从Range序列中滤掉所有奇数只保留偶数IObservableint xs Observable.Range(0, 10); // The numbers 0-9 IObservableint evenNumbers xs.Where(i i % 2 0); evenNumbers.Dump(Where);输出Where--0 Where--2 Where--4 Where--6 Where--8 Where completed语言支持C# 查询表达式Where是所有 LINQ 提供者如 LINQ to Objects 的IEnumerableT都具备的标准操作符。由于 Rx 实现了这些标准操作符C# 查询表达式语法query expression syntax也自然可用——下面的写法与前面的方法调用编译后基本等价IObservableint evenNumbers from i in xs where i % 2 0 select i;原文档的示例大多使用扩展方法而非查询表达式原因有二一是 Rx 实现的部分操作符没有对应的查询语法二是方法调用形式有时更能直观展示发生的过程。惰性订阅与热/冷性质与大多数 Rx 操作符一样Where不会在调用时立即订阅其源。这与 LINQ to Objects 的IEnumerableT版本行为一致只有当你Subscribe由Where返回的IObservableT时它才会转而Subscribe其源并且每次Subscribe都会触发一次对源的订阅。当操作符链式组合时一次对最终序列的Subscribe会沿链一路级联触发对上游的多次Subscribe。这一级联订阅的副作用是Where以及本章大部分操作符既不天生是热的hot也不天生是冷的cold——它只是订阅源源是热它就热源是冷它就冷。谓词调用与终止通知透传从内部机制看当你订阅Where时它会创建自己的IObserverT传给source.Subscribe该 observer 在每次收到OnNext时调用predicate只有谓词返回trueWhere创建的 observer 才会把该元素通过OnNext转发给你传入的 observer。关键在于Where总是把最终的OnCompleted或OnError原样透传。因此即便写出下面这种过滤掉一切的代码IObservableint dropEverything xs.Where(_ false);它虽然会因谓词恒为false而丢弃全部元素但绝不会丢弃错误或完成通知。若你确实只想要源何时完成或失败的信号有更简洁的专门操作符——见下文IgnoreElements。从源码看Where的公开重载定义于 Rx.NET/Source/src/System.Reactive/Linq/Observable.StandardSequenceOperators.csWhereTSource(this IObservableTSource source, FuncTSource, bool predicate)与带索引的FuncTSource, int, bool重载并且对source/predicate为null时抛出ArgumentNullException做了显式校验。IgnoreElements只保留终止通知IgnoreElements扩展方法让你只接收OnCompleted或OnError通知等价于使用谓词恒为false的WhereIObservableint xs Observable.Range(1, 3); IObservableint dropEverything xs.IgnoreElements(); xs.Dump(Unfiltered); dropEverything.Dump(IgnoreElements);输出显示xs产生 1 到 3 后完成而经过IgnoreElements后我们只看到OnCompletedUnfiltered--1 Unfiltered--2 Unfiltered--3 Unfiltered completed IgnoreElements completed其实现定义于 Rx.NET/Source/src/System.Reactive/Linq/Observable.Single.cs。仓库配套测试 IgnoreElementsTest.cs 通过TestScheduler验证了忽略全部元素、仅在OnCompleted时输出、订阅贯穿源生命周期等行为例如IgnoreElements_Basic测试断言即使源产生了 9 个值结果消息集也为空而订阅在整个时间段内保持。OfType按类型过滤某些源会产生多种类型的元素。原文档以 AISAutomatic Identification System船舶自动识别系统为例远洋船舶通过 AIS 广播位置、航向、速度与船名等信息且船名与航行信息的广播频率差异很大。开源 Ais.Net 项目中的ReceiverHost类以IObservableIAisMessage形式暴露消息流其中IAisMessage只报告船舶唯一标识而IVesselNavigation接口报告位置、速度、航向IVesselName接口报告船名。若我们只关心船舶位置而忽略船名很自然会尝试用Where实现// Wont compile! IObservableIVesselNavigation vesselMovements receiverHost.Messages.Where(m m is IVesselNavigation);但这无法编译错误信息为Cannot implicitly convert type System.IObservableAis.Net.Models.Abstractions.IAisMessage to System.IObservableAis.Net.Models.Abstractions.IVesselNavigation原因正如前文所述Where的返回类型永远与输入相同。receiverHost.Messages是IObservableIAisMessageWhere返回的也就是IObservableIAisMessage。虽然谓词保证只有实现IVesselNavigation的消息能通过但 C# 编译器无从理解谓词与输出类型之间的这种关系。为此 Rx 提供了专门的操作符OfType它只保留指定类型的元素——元素必须是指定类型本身、继承自它对于接口则是实现它。于是上面的例子可以这样修复IObservableIVesselNavigation vesselMovements receiverHost.Messages.OfTypeIVesselNavigation();其签名定义于 Rx.NET/Source/src/System.Reactive/Linq/Observable.StandardSequenceOperators.csOfTypeTResult(this IObservableobject? source)配套测试位于 OfTypeTest.cs。位置过滤Positional Filtering有时我们关心的是元素在序列中的位置而非内容本身Rx 为此定义了一组操作符。FirstAsync 与 FirstOrDefaultAsyncLINQ 提供者通常有返回首元素的First。但 Rx 与数据静止型提供者LINQ to Objects、Entity Framework Core不同静止数据的第一项读一下就行而 Rx 源在任意时刻产生数据我们无法预知第一项何时到达。因此 Rx 一般使用FirstAsync它返回一个IObservableT产出源的第一个值后立即完成。Rx 也提供更传统的First方法但它是阻塞式的可能引发问题详见后文阻塞版本小节。结合前面 AIS 的例子报告特定船舶MMSI 为235009890的 HMS Example第一次正在移动的消息。MMSI 即 Maritime Mobile Service Identity船舶移动业务标识。该查询同时用到了本章多个过滤操作符uint exampleMmsi 235009890; IObservableIVesselNavigation moving receiverHost.Messages .Where(v v.Mmsi exampleMmsi) .OfTypeIVesselNavigation() .Where(vn vn.SpeedOverGround 1f) .FirstAsync();流程是先用Where筛出目标船舶的消息再用OfType只看航行消息再用第二个Where忽略未在移动的消息最后由FirstAsync只取第一个正在移动的事件。一旦该船移动moving源立即发出一个IVesselNavigation事件并随即完成。由于FirstAsync可选地接受谓词可以将最后的Where与FirstAsync合并IObservableIVesselNavigation moving receiverHost.Messages .Where(v v.Mmsi exampleMmsi) .OfTypeIVesselNavigation() .FirstAsync(vn vn.SpeedOverGround 1f);空序列行为若输入在产生任何元素前就完成FirstAsync会调用订阅者的OnError抛出带序列不包含元素信息的InvalidOperationException使用谓词重载时若无任何元素匹配谓词也同理。这与 LINQ to Objects 的First一致。前述示例中 AIS 源在应用运行期间会持续广播因此不会出现这种情况。有时我们希望容忍事件缺席。多数 LINQ 提供者还提供FirstOrDefault。Rx 中可结合TakeUntil见下文SkipUntil 与 TakeUntil引入截止时间——例如愿意等待 5 分钟超时即放弃由于可能等不到船移动就完成这里改用FirstOrDefaultAsyncIObservableIVesselNavigation? moving receiverHost.Messages .Where(v v.Mmsi exampleMmsi) .OfTypeIVesselNavigation() .TakeUntil(DateTimeOffset.Now.AddMinutes(5)) .FirstOrDefaultAsync(vn vn.SpeedOverGround 1f);5 分钟后若仍未看到该船以 1 节以上速度移动TakeUntil会退订上游并对FirstOrDefaultAsync传入的 observer 调用OnCompleted。此时FirstOrDefaultAsync与FirstAsync的处理不同它会产出元素类型的默认值IVesselNavigation是接口默认值为null传给订阅者的OnNext后再OnCompleted。简言之这个moving序列始终恰好产出一个元素要么是表示船已移动的IVesselNavigation要么是表示5 分钟内未发生的null。用null表达没发生虽然可行但消费方必须自行区分通知究竟代表目标事件还是事件缺席。Rx 提供了更直接的表示缺席的方式——空序列。你可以想象一个first or empty操作符它要么产出第一个元素要么什么都不产出。这对返回T的 LINQ to Objects 毫无意义但对返回IObservableT的 Rx 是可行的。Rx 内建没有这样的操作符但用更通用的Take可以达到完全相同的效果。TakeTake是标准 LINQ 操作符取序列开头的前几个元素丢弃其余。某种意义上Take是First的推广Take(1)只返回第一项。但二者对缺失元素的响应不同——First/FirstAsync坚持至少要有一个元素空序列会抛InvalidOperationException即便更宽松的FirstOrDefault也坚持产出点什么。Take则不同若输入在达到指定数量前就完成它不会抱怨只是转发源已提供的元素若源只调用了OnCompletedTake就只转发OnCompleted。用Take(5)而源只产了 3 个元素就完成Take(5)会把 3 个元素转发给订阅者后完成。因此可以用Take实现前面设想的FirstOrEmptypublic static IObservableT FirstOrEmptyT(this IObservableT src) src.Take(1);再次强调大多数 Rx 操作符含本章全部操作符本身不热不冷它们跟随源。AIS 的receiverHost.Messages是热的实时广播派生序列因此也是热的。这引出了原文档著名的谐音梗示例IObservableIAisMessage hotTake receiverHost.Messages.Take(1);FirstAsync与Take都从序列开头工作若只关心序列尾部呢LastAsync、LastOrDefaultAsync 与 PublishLastLINQ 提供者通常还提供Last/LastOrDefault与First/FirstOrDefault几乎一样只是返回最后一个元素。同理Rx 以LastAsync/LastOrDefaultAsync提供异步形态也提供可能出问题的阻塞版Last见后文。此外还有 PublishLast语义与LastAsync相近但多订阅的处理方式不同。每次订阅LastAsync返回的序列都会订阅底层源而PublishLast只对底层源发起一次Subscribe。为控制订阅时机PublishLast返回IConnectableObservableT如第 2 章热源与冷源所述它提供Connect方法调用时才会订阅底层源。当这唯一一次订阅收到源的OnComplete后最终值会被投递给所有订阅者PublishLast会记住最终值此后新订阅者一订阅就立即收到该值紧接着发出OnCompleted。它是基于 Multicast 的操作符家族成员后者在后文详述。LastAsync与LastOrDefaultAsync的差异与FirstAsync/FirstOrDefaultAsync相同源空完成时LastAsync报错LastOrDefaultAsync发出默认值后完成。PublishLast对空源又不同源空完成时它产出的序列同样什么都不产出——既无错误也无默认值。关键挑战报告序列末元素比首元素难得多。收到一个元素时无法立即知道它是不是最后一个只有源随后调用OnCompleted才确认。而OnCompleted未必紧随末元素到来如前例用TakeUntil(DateTimeOffset.Now.AddMinutes(5))终止序列末元素与真正完成之间可能隔几分钟。因此LastAsync/LastOrDefaultAsync在源完成前不会发出任何东西这可能导致已从源收到末元素、却迟迟不能转发给订阅者的明显延迟——这是使用它们必须接受的代价。TakeLast与Take相对TakeLast转发序列尾部的元素。例如TakeLast(3)要源的末 3 个元素与Take一样它对元素不足很宽容——源不足 3 个元素时TakeLast(3)直接转发整个序列。TakeLast面临与Last相同的困境不知道何时接近尾部因此必须持有最近看到值的副本内存开销与指定数量成正比。写TakeLast(1_000_000)就要分配容纳 100 万个值的缓冲区——因为在源完成或超过该数量之前它无法确定第一个收到的元素是否属于最终的百万之列。源完成时TakeLast才确定末 N 个元素并逐个通过OnNext转给订阅者。Skip 与 SkipLast若想要Take/TakeLast的严格反操作呢比如某传感器前几次读数总是垃圾值希望忽略前 5 次、等它稳定后再监听——这正是Skip(5)IObservablefloat settled sensor.Readings.Skip(5);SkipLast在序列尾部做同样的事省略末尾指定数量的元素。与前面几个操作符一样它无法预知序列何时结束只有等源发出全部元素并OnComplete后才能确定末 N 个是谁。因此SkipLast会引入延迟使用SkipLast(4)时源要产生第 5 个元素后它才敢转发第 1 个。它不必等OnCompleted/OnError才能行动只要确定某元素不属于待丢弃之列即可开始转发。SingleAsync 与 SingleOrDefaultAsyncLINQ 的Single用于源应恰好包含一个元素的场景多于一个或为空都是错误。Rx 提供SingleAsync返回的IObservableT要么恰好调用一次 observer 的OnNext要么调用OnError表示源出错或未恰好产出一个元素。源为空报错与FirstAsync/LastAsync相同源含多个元素也报错SingleOrDefault则容忍空输入此时产生元素类型默认值。Single/SingleAsync与Last/LastAsync一样收到元素时无法立即确定是否该输出。乍看奇怪既然要求源只提供一个元素收到的第一个不就该是输出吗确实如此——但收到第一项时还不知道源会不会再产生第二个。它不能转发第一项除非源完成以确认不会有更多。可以说SingleAsync的任务是先验证源恰好包含一个元素再转发它失败场景下收到第二个元素时可立即OnError成功场景却必须等源完成才能确认一切正常、发出结果。阻塞版本First / Last / Single[OrDefault]前文多个操作符以Async结尾这略显奇怪——通常 .NET 中以Async结尾的方法返回Task/TaskT而它们返回IObservableT且它们对应的标准 LINQ 操作符并不带Async后缀Entity Framework Core 的Async版本更不同——它们返回TaskT产单个值而非序列。这源于 Rx 早期设计的一个不幸命名决策如果从零设计这些操作符就该叫First、FirstOrDefault等之所以带Async是因为这些操作符在 Rx 2.0 才加入而 Rx 1.0 已占用了原名。下面的代码就用的是Firstint v Observable.Range(1, 10).First(); Console.WriteLine(v);输出1。注意变量v的类型是int而非IObservableint。若用在不会立即产值的 Rx 源上long v Observable.Timer(TimeSpan.FromSeconds(2)).First(); Console.WriteLine(v);运行后会发现First调用直到产出值才返回——它是阻塞操作符。Rx 通常应避免阻塞操作符因为极易造成死锁Rx 的初衷是响应事件的代码坐着干等某个源出值并不符合其精神若确有这种需求往往有更好的建模方式或者 Rx 本就不适合你的场景。如果确实要等待值更好的做法是用Async形态配合 C# 的async/await与 Rx 内置的集成支持long v await Observable.Timer(TimeSpan.FromSeconds(2)).FirstAsync(); Console.WriteLine(v);逻辑效果相同但await不会在等待期间阻塞调用线程降低死锁概率。既然能await就解释了为什么这些方法以Async结尾——虽然它们返回IObservableT而非TaskT具体机制详见《Leaving Rxs World离开 Rx 世界》一章Rx 提供了GetAwaiter扩展方法定义于 Rx.NET/Source/src/System.Reactive/Linq/Observable.Awaiter.csawait一个可观察序列时await会在源完成后结束并返回源产生的最后一个值——这正好适用于FirstAsync/LastAsync这类只产一个元素的场景。不过也存在值立即可得的特例第 3 章的BehaviourSubject恒持有当前值因此First不会真正阻塞——它订阅BehaviourSubjectT而Subscribe在返回前就调用了订阅者的OnNextFirst立即拿到值。当然若用带谓词的重载且当前值不满足谓词First仍会阻塞。ElementAt标准 LINQ 还有按位置取单个元素的操作符ElementAt传入序列中的位置序号在数据静止型提供者中等价于按下标访问数组。Rx 也实现了它ElementAtTSource(this IObservableTSource source, int index)定义于 Rx.NET/Source/src/System.Reactive/Linq/Observable.Aggregates.cs但与First/Last/Single不同Rx不提供ElementAt的阻塞形式。由于任何IObservableT都可被await你可以直接写IAisMessage fourth await receiverHost.Messages.ElementAt(4);若源只产生 5 个值而我们请求ElementAt(5)源完成时ElementAt返回的序列会向订阅者报告ArgumentOutOfRangeException。有三种处理方式优雅处理OnError用.Skip(5).Take(1)忽略前 5 个值只取第 6 个若序列不足 6 个元素得到空序列而不报错使用ElementAtOrDefault索引越界时推送default(T)值目前不提供自定义默认值的能力。时间过滤Temporal FilteringTake/TakeLast用元素数量定义截断点Skip/SkipLast同理。但若想按时刻而非计数来界定范围呢前面已见过一例用TakeUntil把无尽序列变成 5 分钟后完成的序列。这是一族操作符。SkipWhile 与 TakeWhile回到传感器前几次读数不准的场景燃气监测传感器常需把探测组件加热到工作温度才能给出准确读数。前面用Skip(5)很粗糙——怎么知道 5 次足够也许更早就绪了呢真正想要的是丢弃读数直到确定读数有效这正是SkipWhile的用武之地。假设气体传感器同时上报浓度与传感器板温度我们可以直接表达真实需求const int MinimumSensorTemperature 74; IObservableSensorReading readings sensor.RawReadings .SkipUntil(r r.SensorTemperature MinimumSensorTemperature);注意原文档此处实际演示的是SkipWhile的语义在谓词为真期间持续丢弃示例代码中的SkipUntil(r ...)在 Rx 中并不存在该签名重载。Rx 中对应的是SkipWhile(r r.SensorTemperature MinimumSensorTemperature)。若你的 Rx 版本与本文所在仓库一致Rx.NET请使用SkipWhile表达温度达到阈值前一直丢弃的逻辑。下面完整展示SkipWhile的行为。SkipWhile会滤除所有使谓词为真的值直到某个值使谓词为假之后剩余的整个序列原样输出var subject new Subjectint(); subject .SkipWhile(i i 4) .Subscribe(Console.WriteLine, () Console.WriteLine(Completed)); subject.OnNext(1); subject.OnNext(2); subject.OnNext(3); subject.OnNext(4); subject.OnNext(3); subject.OnNext(2); subject.OnNext(1); subject.OnNext(0); subject.OnCompleted();输出4 3 2 1 0 CompletedTakeWhile则相反只要谓词为真就返回所有值第一个使谓词为假的值出现时序列完成var subject new Subjectint(); subject .TakeWhile(i i 4) .Subscribe(Console.WriteLine, () Console.WriteLine(Completed)); subject.OnNext(1); subject.OnNext(2); subject.OnNext(3); subject.OnNext(4); subject.OnNext(3); subject.OnNext(2); subject.OnNext(1); subject.OnNext(0); subject.OnCompleted();输出1 2 3 CompletedSkipWhile/TakeWhile的公开重载含带索引形式定义于 Observable.StandardSequenceOperators.cs。SkipUntil 与 TakeUntil除了SkipWhile/TakeWhileRx 还有SkipUntil/TakeUntil。名字听起来像是同一思想的另一种表达——你可能预期SkipUntil与SkipWhile几乎相同只是前者在谓词返回false期间运行。确实存在这样一对谓词版重载若仅此而已它们就不值得单列。真正有趣的是它们还有用其他方式触发的重载能完成SkipWhile/TakeWhile做不到的事。前面FirstAsync/FirstOrDefaultAsync小节已用过TakeUntil的一个重载接受DateTimeOffset。它包裹任意IObservableT转发源的全部元素直到指定时刻随后立即完成并退订底层源。TakeWhile做不到这一点因为它只在源产出元素时才查询谓词要让源在特定时刻完成TakeWhile只能寄希望于源恰好在那时产出一个元素——它只能因源产出元素而完成。而TakeUntil可以异步完成若指定了 5 分钟后的时刻即使源在那时完全空闲TakeUntil也会完成这依赖 Schedulers调度器机制。TakeUntil还提供接受第二个IObservableT的重载转发源元素直到第二个序列产生一个值无需预先知道那一刻何时到来。SkipUntil有类似的重载由第二个IObservableT决定何时开始转发源元素。这两个重载定义于 Observable.Multiple.cs。重要注意这两个重载要求第二个序列产生一个值来触发开始或结束。若第二个序列未产生任何通知就完成则它毫无作用——TakeUntil会无限期地继续取元素SkipUntil永远不会产出任何东西。换言之这些操作符会把Observable.EmptyT()视同Observable.NeverT()。此外TakeUntil还有接受CancellationToken的重载同样在 Observable.Multiple.cs转发源的元素直到源自身完成或令牌触发取消此时TakeUntil完成。带时间的SkipUntil/TakeUntil重载含可指定IScheduler的版本定义于 Observable.Time.cs。Distinct 与 DistinctUntilChangedDistinct是又一个标准 LINQ 操作符移除序列中的重复项。为此它必须记住源产生过的所有值才能滤掉见过的项。下面的示例用Distinct展示 AIS 消息中首次出现的船舶标识IObservableuint newIds receiverHost.Messages .Select(m m.Mmsi) .Distinct(); newIds.Subscribe(id Console.WriteLine($New vessel: {id}));这里提前用到了Select详见《Transformation of Sequences序列变换》一章用它提取 MMSI 船舶标识。若想同时保留消息细节呢Select提取 id 的做法挡住了这个需求。好在Distinct提供自定义唯一性判定的重载传入一个函数选取任意特征Distinct不再直接比较元素而是比较回调返回值与历史值仅当回调结果是新的才放行。例如IObservableIAisMessage newVesselMessages receiverHost.Messages.Distinct(m m.Mmsi);此时Distinct的输入是IObservableIAisMessage而非上一个示例经Select得到的IObservableuint每次源发出消息就把它传给回调再拿回调返回值与历史值比较仅当 MMSI 前所未见才放行整条消息。效果相似但输出保留完整原始消息——Distinct的输出类型是IObservableIAisMessage。Distinct的重载含IEqualityComparerTSource、FuncTSource, TKey键选择器及其与比较器的组合定义于 Observable.StandardSequenceOperators.cs。除标准Distinct外Rx 还提供DistinctUntilChanged只放行发生变化的通知即仅滤除相邻重复。例如序列1,2,2,3,4,4,5,4,3,3,2,1,1会产出1,2,3,4,5,4,3,2,1。与Distinct记住所有历史值不同DistinctUntilChanged只记住最近发出的一个值新值仅在与该最近值相同时才被滤除。其重载含键选择器与比较器定义于 Observable.Single.cs。下面的示例用DistinctUntilChanged检测某船NavigationStatus航行状态的变化uint exampleMmsi 235009890; IObservableIAisMessageType1to3 statusChanges receiverHost.Messages .Where(v v.Mmsi exampleMmsi) .OfTypeIAisMessageType1to3() .DistinctUntilChanged(m m.NavigationStatus) .Skip(1);例如该船反复报告AtAnchor锚泊状态时DistinctUntilChanged会因状态与上次相同而丢弃每条此类消息一旦状态变为UnderwayUsingEngine使用引擎航行第一个报告该状态的消息就会被放行此后直到状态再次变化回到AtAnchor或变为Moored等才继续放行。末尾的Skip(1)是因为DistinctUntilChanged总会放行它看到的第一个消息——我们无法知道它是否真的代表状态变化而船舶每几分钟报告一次状态、却远不常改变状态因此第一个报告大概率不代表变化丢掉这一条确保statusChanges只在我们确定状态确实变了时才发出通知。小结以上就是 Rx 过滤操作符的快速巡礼。它们相对简单但正如我们已经看到的Rx 的力量在于操作符的可组合性——Where、OfType、TakeUntil、DistinctUntilChanged等可以像积木一样自由拼接如 AIS 示例中的多级组合针对同一数据流表达极其精准的筛选条件。过滤操作符是应对这个信息过载时代数据洪流的第一道防线现在你已掌握如何用各种准则去除数据下一步即可进入《Transformation变换》一章学习对保留下来的数据进行转换加工。若需在本地验证以上行为可查看仓库中的 System.Reactive 源码目录公开操作符入口与 Tests.System.Reactive 测试目录含IgnoreElementsTest.cs、OfTypeTest.cs等覆盖边界行为的用例并参考 Rx.NET 项目主文档了解IObservableT与热/冷源基础概念。赞分享后端【免费下载链接】reactiveThe Reactive Extensions for .NET项目地址https://gitcode.com/gh_mirrors/re/reactive点击查看免费下载相关推荐Freewall实战5种惊艳布局效果让你的网站脱颖而出Freewall实战5种惊艳布局效果让你的网站脱颖而出 想要为你的网站创建令人惊艳的动态网格布局吗Freewall是一个强大而灵活的响应式网格布局引擎可以UI库/组件RxJS 4 的 filter 与 where 操作符基于谓词过滤 Observable 序列的完整指南RxJS 4 的 filter 与 where 操作符基于谓词过滤 Observable 序列的完整指南 导读 filter 别名 where 是 Rea后端League Akari开发者指南基于LCU API的插件开发与功能扩展League Akari开发者指南基于LCU API的插件开发与功能扩展 League Akari是一款强大的《英雄联盟》客户端一体化工具包通过LCU AP桌面应用上一篇V-Calendar终极国际化配置指南快速实现多语言和本地化支持下一篇终极自定义滚动条指南Baron让你的网页滚动更优雅 创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表