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
为什么Stream除了Value还有Time与Lifecy…
Collection是已经存在的一组值, IEnumerable<T> 由consumer拉取, IObservable<T> 则由source随时间推送。Observable contract包含零到多个OnNext,随后至多一个OnError或OnCompleted;su…
- 2
Representing data streams wit…
IObservable<T> 只有Subscribe入口,observer提供三个notifications。它不在type中承诺backpressure、thread、scheduler、cold/hot、replay或multicast,这些都要由source/operator co…
- 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吗。答案都是否。
↡由OnNext序列、唯一terminal signal与subscription disposal共同定义的push stream生命周期契约。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。
↡每次Subscribe是否重新执行source以及多个subscribers是否共享同一production过程的语义选择。切换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。
↡Operator对多个source的subscription、notification ordering、completion与cancellation所建立的组合语义。切换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。
↡以Scan、Window或显式state machine在多个notifications和时间边界之间保持受控状态的逻辑。切换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<T>适合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
- Observable不仅是values,还包含terminal signals、subscription ownership和scheduler。
- Cold/hot、unicast/multicast与replay必须显式说明,type本身不承诺。
- Operators编码ordering、completion和cancellation,需用virtual-time timeline验证。
- 跨事件逻辑应提取pure temporal state transition,并定义event/processing time。
- 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?应写成可组合的输入—输出契约,并把环境读取、状态改变与失败显式放在边界。用确定输入运行正常、空值和失败样本,同时记录返回值与副作用轨迹,证明结论不依赖隐藏状态。