Chapter 14. Data streams and the Reactive Extensions

覆盖IObservable contract、source creation、Rx operators、cross-event state与适用边界,把time、termination、subscription和resource ownership显式化。

学习目标

  • 能解释IObservable的OnNext/OnError/OnCompleted与Dispose契约,区分cold/hot和unicast/multicast source
  • 能实现resource-safe observable factory,并设计Select、Where、Merge、Switch等operator的ordering/cancellation tests
  • 能分析跨多事件的window、Scan、timeout和retry状态,判断何时选择IObservable、IAsyncEnumerable或Channel

机制总览

Chapter 14. Data streams and the Reactive Extensions:机制路径

  1. 1

    为什么Stream除了Value还有Time与Lifecy…

    Collection是已经存在的一组值, IEnumerable<T> 由consumer拉取, IObservable<T> 则由source随时间推送。Observable contract包含零到多个OnNext,随后至多一个OnError或OnCompleted;su…

  2. 2

    Representing data streams wit…

    IObservable<T> 只有Subscribe入口,observer提供三个notifications。它不在type中承诺backpressure、thread、scheduler、cold/hot、replay或multicast,这些都要由source/operator co…

  3. 3

    Creating IObservables

    优先使用 Observable.Return/Empty/Throw/Defer/FromAsync/Using 等已验证constructors。自定义 Observable.Create 时,subscribe function必须返回disposable,负责解绑event、取消token/timer并防止terminal后继续emit。

先按顺序建立机制,再进入实验切换阶段并检查失效证据。

章级决策实验

Chapter 14. Data streams and the Reactive Extensions:机制与证据

切换《Chapter 14. Data streams and the Reactive Extensions》的三个关键教学阶段,先解释机制,再用运行与失败证据验证结论。

选择推理阶段

当前阶段 · 为什么Stream除了Value还有Time与Lifecy…

Collection是已经存在的一组值, IEnumerable<T> 由consumer拉取, IObservable<T> 则由source随时间推送。Observable contract包含零到多个OnNext,随后至多一个OnError或OnCompleted;su…

可核验证据

以确定输入重复运行「为什么Stream除了Value还有Time与Lifecy…」的最小管线,用属性测试、状态快照和副作用调用轨迹核对返回值、失败传播与资源边界。

学完《Chapter 14. Data streams and the Reactive Extensions》后,应能从输入和前置条件推导状态变化,并用可重复的构建、运行或边界测试证明结果。

失效—证据矩阵

Chapter 14. Data streams and the Reactive Extensions:失效与核验

为什么Stream除了Value还有Time与Lifecy…

典型失效

若把「为什么Stream除了Value还有Time与Lifecy…」只写成函数式术语而不隔离副作用、状态和失败分支,组合后的程序仍会依赖隐藏时序,无法从输入稳定推导结果。

核验证据

以确定输入重复运行「为什么Stream除了Value还有Time与Lifecy…」的最小管线,用属性测试、状态快照和副作用调用轨迹核对返回值、失败传播与资源边界。

Representing data streams wit…

典型失效

若把「Representing data streams wit…」只写成函数式术语而不隔离副作用、状态和失败分支,组合后的程序仍会依赖隐藏时序,无法从输入稳定推导结果。

核验证据

以确定输入重复运行「Representing data streams wit…」的最小管线,用属性测试、状态快照和副作用调用轨迹核对返回值、失败传播与资源边界。

Creating IObservables

典型失效

若把「Creating IObservables」只写成函数式术语而不隔离副作用、状态和失败分支,组合后的程序仍会依赖隐藏时序,无法从输入稳定推导结果。

核验证据

以确定输入重复运行「Creating IObservables」的最小管线,用属性测试、状态快照和副作用调用轨迹核对返回值、失败传播与资源边界。

每个判断都必须能落到观测、测试或产物,不能只凭代码表面推测。

为什么Stream除了Value还有Time与Lifecycle

Collection是已经存在的一组值,IEnumerable<T>由consumer拉取,IObservable<T>则由source随时间推送。Observable contract包含零到多个OnNext,随后至多一个OnError或OnCompleted;subscription dispose是consumer停止观察并释放ownership,不等于source成功完成。忽略terminal和dispose,只关注values,会留下timer、socket和event handler泄露。

先预测:OnError后还能OnNext吗;每个Subscribe一定创建独立source吗;Dispose必然触发OnCompleted吗;Merge保持两个sources全局顺序吗;Retry只重新处理失败item吗。答案都是否。

Representing data streams with IObservable

IObservable<T>只有Subscribe入口,observer提供三个notifications。它不在type中承诺backpressure、thread、scheduler、cold/hot、replay或multicast,这些都要由source/operator contract补充。Cold source通常每个subscriber独立执行recipe;hot source与subscriber无关地生产,late subscriber可能错过过去values。

IDisposable subscription = prices.Subscribe(
    onNext: price => view.Update(price),
    onError: error => view.ShowFailure(error),
    onCompleted: () => view.MarkComplete());
 
lifetime.Add(subscription);

Notification serialization也需明确:标准Rx operators通常保证一个subscription的callbacks不并发,但手写source可能从多个threads调用observer。必须用scheduler/Serialize或single owner保护contract。Observer callback抛异常不等同source OnError;UI adapter应在边界处理自身fault,不让一个consumer破坏shared source。

分步1 / 3

切换OnNext、OnError、OnCompleted与Dispose

Creating IObservables

优先使用Observable.Return/Empty/Throw/Defer/FromAsync/Using等已验证constructors。自定义Observable.Create时,subscribe function必须返回disposable,负责解绑event、取消token/timer并防止terminal后继续emit。Defer把source construction放到每次Subscribe;若直接在外部创建Task,再用FromAsync包装,Task可能在subscriber出现前已经启动。

IObservable<Reading> Readings(ISensor sensor) =>
    Observable.Create<Reading>(observer =>
    {
        void OnReading(object? _, Reading value) => observer.OnNext(value);
        sensor.Reading += OnReading;
        sensor.Start();
 
        return Disposable.Create(() =>
        {
            sensor.Reading -= OnReading;
            sensor.Stop();
        });
    });

Factory需定义一个subscriber dispose是否停止shared source,reference count归零才停止,还是source由application lifetime拥有。Publish().RefCount()把cold变shared hot,但disconnect/reconnect可能重置sequence;Replay(1)会缓存last value,也会延长其对象图生命周期。测试first/late/two subscribers、last unsubscribe、source fault和resubscribe。

Scheduling and thread-affinity evidence

SubscribeOn影响subscription/unsubscription side effects发生在哪里,ObserveOn影响后续observer notifications在哪个scheduler。两者不可互换;也不要在library内部随意硬编码UI scheduler。Scheduler是dependency,可在tests使用virtual time精确推进timer/debounce,而不是sleep等待造成慢且不稳定的测试。

生产stream要记录source timestamp、ingest timestamp和processing lag,区分event time与arrival time。跨thread UI更新必须在明确ObserveOn boundary切换;CPU-heavy Select不应阻塞event loop。若需要并发处理,SelectMany/merge还要有capacity和ordering策略,Rx本身不提供无限吞吐。

Transforming and combining data streams

Select逐值映射,Where过滤,DistinctUntilChanged抑制相邻重复;Merge交错sources,Concat保持先后subscription,CombineLatest用每个source最新值,Zip按index配对,Switch只观察最新inner source并unsubscribe旧inner。选择operator要从domain timing问题出发,不按名字试到“看起来能跑”。

IObservable<SearchResult> results = queries
    .Select(text => text.Trim())
    .Where(text => text.Length >= 2)
    .DistinctUntilChanged()
    .Throttle(TimeSpan.FromMilliseconds(250), scheduler)
    .Select(query => api.Search(query).ToObservable())
    .Switch();

上述Switch适合“只要最新query”,旧request结果不会更新UI;但unsubscribe能否真正取消HTTP取决于adapter token wiring。Merge适合所有inner results都重要,Concat适合顺序effects。Combining sources时写出completion rule、one source errors时是否终止、empty source行为及cross-source ordering。

分步1 / 3

切换Select、Where、Merge与Switch

Implementing logic that spans multiple events

跨事件逻辑需要state和time。Scan是stream上的fold,每个event产生new accumulator并可输出intermediate state;Buffer/Window按count/time分组;Pairwise比较相邻值;Timeout把silence变为terminal/fallback。复杂CEP先写state machine:输入event、current state、timer,输出new state/actions,再由Rx连接sources和scheduler。

例如登录失败告警不是简单Buffer(3):规则可能要求同账号、五分钟窗口、成功后清零、锁定后忽略、late event按event time处理。把这些写成immutable accumulator与pure Scan transition,virtual-time tests覆盖边界tick、out-of-order和clock skew。若state需跨process持久化,Rx in-memory Scan不够,需要durable stream processor/checkpoint。

分步1 / 3

切换Buffer、Scan、Timeout与Retry

Backpressure and overload boundaries

IObservable是push contract,没有built-in demand signal。Producer快于consumer时,选择drop/sample/latest、buffer with bound、batch、pause adapter或转Channel。无限Buffer只是把延迟变成memory failure;ObserveOn也会建立queue。Metrics必须包含queue depth、oldest age、drop count和processing rate,alert按SLO而非OOM后再发现。

若每个item处理异步且昂贵,SelectMany无限merge会增加in-flight。采用bounded concurrency operator/Channel worker,并明确output order。Market data UI可sample latest;审计events不能drop,应使用durable broker与consumer offsets。Backpressure policy来自业务价值,而非统一library default。

When should you use IObservable?

IObservable适合多值push、event/UI/sensor流、时间operators和多source组合,尤其需要virtual-time reasoning。IAsyncEnumerable&lt;T&gt;适合consumer-driven async iteration与自然backpressure;Channel适合producer/consumer queue、bounded capacity与worker ownership;普通Task适合单个未来结果;domain event log适合durable replay。它们可在边界转换,但转换会改变hot/cold、buffering、cancellation和error semantics。

决策表至少比较cardinality、push/pull、subscriber count、backpressure、durability、replay、threading和operator需求。若只返回一个HTTP结果,用Task;若读取分页stream且consumer控制速度,用IAsyncEnumerable;若多个listeners组合UI events与debounce,Rx更自然;若不能丢消息并需跨重启恢复,先选durable broker而非内存Observable。

Production gate: subscriptions must have owners

每个Subscribe都要能回答“谁Dispose、何时Dispose、terminal后还做什么”。UI绑定到view lifetime,service shared subscription绑定host lifetime,per-request subscription绑定request token。Static subject和忘记dispose会保留observer/object graph;定期heap snapshot和subscription metrics可暴露泄漏,但代码review应先要求owner。

Integration test用真实adapter验证dispose确实解绑handler/取消I/O,virtual scheduler验证operator semantics,load test验证queue/backpressure。三个层次分别证明resource、logic与capacity,不能只靠marble diagram宣称production ready。

本章回顾:Reactive正确性来自完整Lifecycle

  1. Observable不仅是values,还包含terminal signals、subscription ownership和scheduler。
  2. Cold/hot、unicast/multicast与replay必须显式说明,type本身不承诺。
  3. Operators编码ordering、completion和cancellation,需用virtual-time timeline验证。
  4. 跨事件逻辑应提取pure temporal state transition,并定义event/processing time。
  5. IObservable无内建backpressure;按业务选择bound/drop/durable queue或其他abstraction。

练习

问题 1:怎样验证一个Event-to-Observable adapter不泄漏?

问题 2:搜索提示为何常用Switch而非Merge?

问题 3:何时应从IObservable改为Channel?

术语表

名词解释

本章出现的专业名词,用大白话再讲一遍。

observable lifecycle contract
cold-hot boundary
subscription resource boundary
stream composition policy
temporal state transition

原版目录概念补充核对

以下条目补齐官方目录中容易被示例主线掩盖的概念。它们不重复罗列目录,而是明确每项概念的机制、适用边界和验收证据。

Creating IObservables:机制、边界与证据

Chapter 14. Data streams and the Reactive Extensions中的Creating IObservables应写成可组合的输入—输出契约,并把环境读取、状态改变与失败显式放在边界。用确定输入运行正常、空值和失败样本,同时记录返回值与副作用轨迹,证明结论不依赖隐藏状态。

Transforming and combining data streams:机制、边界与证据

Chapter 14. Data streams and the Reactive Extensions中的Transforming and combining data streams应写成可组合的输入—输出契约,并把环境读取、状态改变与失败显式放在边界。用确定输入运行正常、空值和失败样本,同时记录返回值与副作用轨迹,证明结论不依赖隐藏状态。

When should you use IObservable?:机制、边界与证据

Chapter 14. Data streams and the Reactive Extensions中的When should you use IObservable?应写成可组合的输入—输出契约,并把环境读取、状态改变与失败显式放在边界。用确定输入运行正常、空值和失败样本,同时记录返回值与副作用轨迹,证明结论不依赖隐藏状态。

讨论

评论区加载中…