本节摘要:点击、消息、行情、传感器读数——这些随时间到达的事件,命令式代码用回调逐个应对,函数式代码把它们看作"时间维度上的集合",于是 map、filter、reduce 全部重新可用。本节建立事件流的集合观,拆解主流响应式库的操作符体系与背压机制,并用一个搜索联想的完整案例对比回调与流两种写法。
场景:搜索框输入联想——用户每敲一个字符,等三百毫秒停顿,发请求,展示结果;发新请求时取消旧的;请求返回的顺序要跟输入对齐,不能旧盖新。命令式回调版本:
// 回调版:五个关注点散落在定时器、标志位与闭包之间 let timer = null; let seq = 0; let lastSeqDone = 0; input.addEventListener("input", (e) => { clearTimeout(timer); timer = setTimeout(() => { const mySeq = ++seq; fetch("/suggest?q=" + encodeURIComponent(e.target.value)) .then((r) => r.json()) .then((list) => { if (mySeq > lastSeqDone) { // 乱序防护:手工比较序号 lastSeqDone = mySeq; render(list); } }); }, 300); });
二十行里,防抖、取消、乱序防护、渲染四种关注点互相纠缠,每加一个需求(失败重试、空值跳过、组件卸载时清理)都要动多处全局状态。问题不在程序员,在模型不对:事件流的"时间结构"(防抖、节流、乱序、合并)没有被建模,只能靠手搓状态机模拟。
函数式的重新建模只有一句话:把事件流(Observable/Stream)当作随时间到达的集合。列表能做的变换,流都能做,只是"遍历"变成了"订阅":
// RxJS 流版:同样的需求,关注点各归其位 import { fromEvent } from "rxjs"; import { debounceTime, map, switchMap, filter, distinctUntilChanged } from "rxjs/operators"; const input$ = fromEvent(input, "input").pipe( map((e) => e.target.value.trim()), // 变换:取值并清洗 filter((q) => q.length > 0), // 筛选:空串不发请求 distinctUntilChanged(), // 去重:内容没变不重发 debounceTime(300), // 时间结构:停顿300毫秒 switchMap((q) => fetch("/suggest?q=" + encodeURIComponent(q)) .then((r) => r.json())), // 切换:新请求自动作废旧请求 ); input$.subscribe(render); // 唯一的副作用出口
逐个操作符对应回调版的手工劳动:debounceTime 取代 timer 手搓,switchMap 取代序号比较(它内部保证"只留最新内层流",旧请求的结果直接丢弃),filter、map 与第 3 章的同名函数语义完全一致。同一套集合词汇,从列表平移到了时间轴上——这就是响应式流的学习成本为什么对函数式老手如此之低。
一个思维校准:流上的 map 作用于"每个事件",而不是"整个流"——input$ 流过 map 后还是流,事件被逐个变换。这正是第 4 章函子的又一次现身:事件流是函子、是应用函子、也是 Monad(switchMap 就是 bind 的流版本:每个事件映射出一个新的子流,自动展平)。
操作符数量庞大(主流库上百个),但按功能归组只有六族,记住六族就能按图索骥:
变换族 map / scan / buffer 逐事件变换 / 携带累积器 / 打包成批 筛选族 filter / debounceTime / throttle / distinct / take 合并族 merge / combineLatest / zip / withLatestFrom 多流协作 展平族 switchMap / mergeMap / concatMap / exhaustMap 流中流 错误族 retry / retryWhen / catch / onErrorReturn 多播族 share / publish / refCount 一条流多个订阅者共享执行
scan 值得单独点名:它是第 3 章 reduce 的流版本——累积器滚过每个事件。点击计数、增量聚合、状态折叠,全部是 scan 的换皮:
// 双击计数器:scan 就是流上的 reduce clicks$.pipe( bufferWhen(() => interval(250)), // 250毫秒内的点击打包 map((batch) => batch.length), filter((n) => n >= 2), ).subscribe((n) => console.log("检测到连击", n));
**背压(backpressure)**是流体系区别于回调的独门机制:消费方跟不上生产方时怎么办?回调世界没有答案(事件堆积或丢失,看实现心情);流体系把它变成显式选择——buffer(排队)、throttle(采样)、latest(只留最新)、onBackpressureDrop(丢弃并记录)。生产速率与消费速率的错配,从"偶发线上事故"变成"管道上一处可见的配置"。
严格说,响应式流(RxJS、Reactor、Akka Streams)与函数式响应式编程(FRP)不是一回事。流是离散的:事件的集合,事件之间没有值;FRP 是连续的:值随时间连续变化(行为/Behavior),任意时刻都可询问"现在的值是多少"。界面上的鼠标位置、滑块数值是典型的连续量——FRP 把"位置随时间"建模为时间到值的函数,而不是"每个鼠标移动事件"的序列。实践影响:离散流适合消息类业务(订单、通知),连续行为适合界面状态推导(禁止提交 = 表单无效 或 请求中)。前端框架中,React 的"UI 是状态的函数"接近 FRP 精神,RxJS 是离散流代表,Vue 的响应式数据介于两者之间。
把案例升级为接近生产的形态,展示六族操作符的协作,并处理副作用收口:
// 搜索联想生产版:取消、错误、共享、清空,各一处 import { fromEvent, EMPTY, merge } from "rxjs"; import { debounceTime, map, switchMap, catchError, share, filter, tap } from "rxjs/operators"; const search$ = fromEvent(input, "input").pipe( map((e) => e.target.value.trim()), debounceTime(300), filter((q) => q.length >= 2), share(), // 多播:清空分支与结果分支共用一条流 switchMap((q) => fetch(`/suggest?q=${encodeURIComponent(q)}`) .then((r) => { if (!r.ok) throw new Error("HTTP " + r.status); return r.json(); }), ), catchError(() => EMPTY), // 失败静默,流不断 ); merge( search$.pipe(tap(render)), // 副作用唯一出口之一 search$.pipe(filter((list) => list.length === 0), tap(clear)), ).subscribe();
判别标准收束:什么时候该上流——事件之间有时间结构(防抖、合并、乱序)或多个事件源要协作时,流的操作符直接表达时间语义,收益最大;什么时候不必上——单点事件单点响应(一个按钮一个请求),回调或 async/await 更直白,引入流是自增复杂度。另外两条工程忠告:流的订阅生命周期要与组件/连接生命周期绑定(忘了退订就是内存泄漏,这是响应式最常见的翻车点);调试流要靠操作符日志(tap(console.log))逐段观测,因为栈轨迹在流里基本失效。
⚠️ 常见坑:
mergeMap与switchMap一字之差、行为天壤——前者所有子流并存(可能请求风暴),后者只留最新。表单联想用 switch,点赞收藏用 merge(每个都要落地),选错即事故。
事件流的组织方案就位,本章还剩最后一笔账:性能。不可变的分配、惰性的驻留、流的调度——下一节把收益与成本两边全部摊开。