- 后端
【免费下载链接】reactive
The Reactive Extensions for .NET
本篇技术指南以 Ix.NET/Documentation/Museum/OldReadme.md(Rx.NET 项目 2023 年前的主 README 存档)为骨架,系统讲解 Reactive Extensions(Rx)的核心抽象:用IObservable<T>表示异步数据流、用 LINQ 运算符查询数据流、用 Scheduler 参数化并发;并结合当前仓库源码(Rx.NET/Source/src/System.Reactive)与最新主 README.md 验证这些概念在实现层面的落点。读完本文,你将理解 Rx 的三大支柱及其在 .NET 生态中的定位,并能在仓库源码中快速定位对应实现。
为什么会有这份 "OldReadme" 存档
按照该文档的说明:到 2023 年初,主 README 已经累积了大量信息,其中大部分对于第一次访问仓库、想了解"这个项目是干什么的"的访客来说,细节层次并不合适。但旧内容中大多数信息对某些人仍有用,于是维护者将较冷僻的部分移入OldReadme.md,以免彻底丢失。也就是说,这份存档记录的是 Rx 项目早期对自身定位的经典表述——它聚焦于"细节",未必直指"Rx 到底用来解决什么问题",但其中的核心模型(Observables + LINQ + Schedulers)至今仍是 Rx.NET 的架构基石。
当前仓库的实际面貌已大幅演进:主 README.md 现在用"live data streams(实时数据流)"来阐释 Rx 的价值,并同时托管四个相关库——Rx.NET、实验性的 AsyncRx.NET、Interactive Extensions(Ix.NET)以及System.Linq.Async。本文接下来以存档文档为主线,逐层展开。
Rx 的经典定义:Observables + LINQ + Schedulers
存档文档给出了 Rx 最经典的概括:
Rx 是一个使用**可观察序列(observable sequences)**和LINQ 风格查询运算符来组合异步与事件驱动程序的库。开发者用
IObservable<T>表示异步数据流,用LINQ 运算符查询异步数据流,用Scheduler参数化异步数据流中的并发。简言之:Rx = Observables + LINQ + Schedulers。
这一定义拆开看有三层含义:
- 表示(Represent):异步数据流的载体是
IObservable<T>,即"可观察序列"接口; - 查询(Query):数据流上可以施加 LINQ 查询运算符(
Where、Select、Aggregate等); - 参数化(Parameterize):并发与调度行为通过 Scheduler 注入,使代码与具体的线程模型解耦。
现实问题:异步与事件编程的复杂性
文档指出,无论传统桌面应用还是 Web 应用,都必须面对异步与事件驱动编程:桌面应用有 I/O 操作和计算密集型任务,可能长时间占用并阻塞其他活动线程;异常处理、取消、同步在这些场景下困难且易错。Rx 的切入点是:把来自不同来源(股票行情、推文、计算机事件、Web 服务请求等)的多条异步数据流统一表示出来,并用IObserver<T>接口订阅事件流——每当事件发生,IObservable<T>就通知已订阅的IObserver<T>。
核心接口:IObservable 与 IObserver
IObservable :数据流的"生产者"
IObservable<T>表示一个可被订阅的数据源,核心能力是Subscribe。从源码结构看,Rx.NET 在其基础上提供了一整套扩展与实现,例如 Rx.NET/Source/src/System.Reactive/AnonymousObservable.cs、AnonymousObserver.cs 等,用于把委托转换成标准的可观察/观察者对象,供Observable.Create之类的工厂方法使用。
IObserver :数据流的"消费者"
IObserver<T>定义了对事件流的三类响应:
OnNext(T value):数据流产生新元素;OnError(Exception error):数据流出错终止;OnCompleted():数据流正常结束。
仓库中 ObserverBase.cs 定义了观察者的基类,并对协议的合法性(如在已完成/出错后不能再OnNext)做了约束,这正是 Rx 把异步事件流"协议化"的关键。
一个贯穿始终的例子
主 README.md 给出了与存档文档一脉相承的示例:假设trades是一个IObservable<Trade>,那么可以用 LINQ 查询只保留成交量超过一百万的交易:
var bigTrades = from trade in trades where trade.Volume > 1_000_000 select trade;查询表达式语法只是方法调用的语法糖,等价于:
var bigTrades = trades.Where(trade => trade.Volume > 1_000_000);随后订阅并消费结果:
bigTrades.Subscribe(t => Console.WriteLine($"{t.Symbol}: trade with volume {t.Volume}"));这里bigTrades同样是IObservable<Trade>,且只在每次trades产生符合条件的事件时触发回调——这就是"查询实时数据流"的含义:结果不是一次性物化的集合,而是一条持续推送的派生流。
用 LINQ 运算符查询数据流:源码实现落点
存档文档强调:因为可观察序列是数据流,可以用 Observable 扩展方法实现的标准 LINQ 运算符来查询它们,从而方便地对多条事件流做过滤、投影、聚合、组合以及基于时间的操作;此外还有一批反应式流专属运算符,能写出更强大的查询;取消、异常和同步也由 Rx 提供的扩展方法优雅处理。
在 Rx.NET/Source/src/System.Reactive/Linq 下可以找到这些运算符的实现,例如:
- Observable.StandardSequenceOperators.cs 集中了
Cast、DefaultIfEmpty、Distinct、Select、Where等标准序列运算符。以Where为例,其签名是Where<TSource>(this IObservable<TSource> source, Func<TSource, bool> predicate),先对source与predicate做空值校验,再交给内部实现s_impl.Where(source, predicate);Select则提供带索引与不带索引两种重载(Func<TSource, TResult>与Func<TSource, int, TResult>); - Observable.Aggregates.cs 对应聚合类运算;
- Observable.Creation.cs、Observable.Concurrency.cs、Observable.Time.cs 分别覆盖序列创建、并发控制与基于时间的操作;
- 多源组合类运算符则分布在 Observable.Multiple.cs、Observable.Multiple.CombineLatest.cs、Observable.Multiple.Zip.cs 等文件中。
值得一提的是,这些运算符文件多采用"公开 API + 内部s_impl实现"的拆分模式(见 Observable_.cs),公开层负责参数校验与文档,实际算法在内部类中实现,便于测试与维护。
Rx 与其他数据模型的关系:Pull/Push 二维表
存档文档用一张经典二维表说明 Rx 与同步数据流、单值异步计算的互补关系:
| 单返回值 | 多返回值 | |
|---|---|---|
| Pull / 同步 / Interactive | T | IEnumerable<T> |
| Push / 异步 / Reactive | Task<T> | IObservable<T> |
这张表揭示了四类数据模型的分工:
T:同步、单值(普通方法返回值);IEnumerable<T>:同步、多值(拉取式迭代);Task<T>:异步、单值(可等待的单个异步结果);IObservable<T>:异步、多值(推送式事件流)。
Rx 的定位正是补齐"异步 + 多值"这一象限,并与另外三种模型顺畅互操作。这一点在仓库中得到印证:System.Reactive/Linq/Observable.Conversions.cs 提供IEnumerable<T>与IObservable<T>之间的转换,Observable.Async.cs 处理与Task<T>的互操作,Observable.Awaiter.cs 则允许直接await一个可观察序列。
演进:AsyncRx.NET 与 IAsyncObservable
主 README.md 补充了这一定位的后续演进:Rx 设计早于 C# 的async/await,其原始设计假定处理通知的代码同步执行,因此存在某些无法使用async的场景。仓库内的 AsyncRx.NET(实验性预览)通过定义IAsyncObservable<T>解除这一限制,允许观察者使用异步代码:
bigTrades.Subscribe(async t => await bigTradeStore.LogTradeAsync(t));对应实现位于 AsyncRx.NET/System.Reactive.Async/IAsyncObservable.cs 与 IAsyncObserver.cs。同时,Ix.NET 中的System.Linq.Async为IAsyncEnumerable<T>提供标准 LINQ 运算符实现,与 Rx 形成互补的完整异步数据流生态。
Schedulers:参数化并发
存档文档将 Scheduler 列为 Rx 三大支柱之一:开发者用 Scheduler参数化异步数据流中的并发。这意味着查询本身不直接绑定到某个线程或调度策略,而是把"在何时、用哪个上下文执行"的决定权交给调度器,从而让代码可测试、可移植。
仓库中 Concurrency/IScheduler.cs 定义了调度器的核心契约,Concurrency 目录下还包含周期性调度(ISchedulerPeriodic)与长任务调度(ISchedulerLongRunning)等扩展接口,以及DefaultScheduler、TaskPoolScheduler、NewThreadScheduler、SynchronizationContextScheduler等具体实现。结合 Observable.Concurrency.cs 中的SubscribeOn/ObserveOn等运算符,可以在订阅与通知两个阶段分别指定调度器,实现"查询逻辑与线程模型解耦"。
Flavors of Rx:跨语言生态
存档文档列出了 Rx 在不同语言/运行时下的实现:
- Rx.NET(即本仓库):面向 .NET 的 Reactive Extensions,基于
IObservable<T>与 LINQ 风格查询运算符; - RxJS:面向 JavaScript(浏览器与 Node.js)的实现;
- RxJava:面向 JVM 的实现;
- RxScala:面向 Scala 的实现;
- RxCpp:以 C 与 C++ 提供实现;
- RxPy:面向 Python 3 的实现。
它们共享同一套"可观察序列 + 查询运算符"的编程模型,只是载体语言不同。需要注意:上述除 Rx.NET 外的项目均不在本仓库内,本文仅转述存档文档的生态介绍;主 README.md 亦提到 RxJS 在 UI 编程中的流行,以及 ReactiveUI 对 Rx 的深度使用,印证了这一模型在跨语言层面的影响力。
应用案例:Tx 与 LINQ2Charts
存档文档还记录了两个基于 Rx 的应用示例:
- Tx:一组展示如何使用"LINQ to events"的代码示例,包括对实时数据的常驻查询(standing queries),以及对过去历史数据的查询(来自 trace 与 log 文件),目标数据源涵盖 ETW、Windows 事件日志与 SQL Server 扩展事件;
- LINQ2Charts:一个 Rx 绑定的示例,类似 LINQ to XML,允许开发者用 LINQ 以简单方式创建/修改/更新图表,从而避免直接处理 XML 等底层数据结构。
这两个示例同样不在本仓库内,属于存档文档对早期生态的记录;它们展示了"把事件流当作可查询序列"这一思想在日志分析、可视化绑定等领域的延伸价值。
仓库现状与进一步阅读
存档文档描述的模型在今天依然是仓库的根脉,但仓库面貌已更新。若想深入,可以按以下路径阅读:
- 最新介绍与四个库的定位:README.md(含交易流示例、AsyncRx/Ix 说明、获取 NuGet 包的渠道与夜间构建 feed);
- Rx.NET 完整源码:Rx.NET/Source/src/System.Reactive(核心接口、运算符、调度器、Subject、Disposable 等约 360 个文件);
- 运算符测试:Rx.NET/Source/tests/Tests.System.Reactive(每个运算符都有对应行为验证);
- 免费电子书《Introduction to Rx.NET 2nd Edition》的 Markdown 全文:Rx.NET/Documentation/IntroToRx;
- Rx.NET 发布历史:Rx.NET/Documentation/ReleaseHistory,Ix 的发布历史见 Ix.NET/Documentation/ReleaseHistory。
小结
通过这份"初代 README"存档,可以清晰看到 Rx 数十年不变的三个核心命题:用IObservable<T>表示异步数据流、用 LINQ 运算符查询数据流、用 Scheduler 参数化并发。仓库源码(Observable.StandardSequenceOperators.cs、Concurrency/IScheduler.cs)与最新主 README 均与这些命题一一对应,而 AsyncRx.NET 与 System.Linq.Async 则把"异步 + 多值"的象限进一步拓宽。理解这三根支柱,是使用 Rx.NET 编写可靠、可组合的异步事件程序的第一步。
- 后端
【免费下载链接】reactive
The Reactive Extensions for .NET
相关推荐
RxJS 4.0 实战指南:用 Observables、Operators 与 Schedulers 组合异步与事件驱动程序
RxJS 4.0 实战指南:用 Observables、Operators 与 Schedulers 组合异步与事件驱动程序 本指南以当前仓库根目录的 read
后端RxJS 4 模块化架构实战:基于 @rxjs/rx 的 Observables、Operators 与 Schedulers 组合式异步编程指南
RxJS 4 模块化架构实战:基于 @rxjs/rx 的 Observables、Operators 与 Schedulers 组合式异步编程指南 本指南以仓库
后端libuv事件循环:异步编程的核心引擎
libuv事件循环:异步编程的核心引擎 本文深入解析了libuv事件循环的11个执行阶段、UV_RUN运行模式的区别、定时器处理机制以及线程安全性与多事件循环的
网络通信异步编程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考