- 后端
【免费下载链接】RxJS
The Reactive Extensions for JavaScript
导读
Rx.Observable.prototype.takeLastWithTime(duration)是 RxJS v4(Reactive Extensions for JavaScript)中一个按时间维度截取序列末尾元素的时间类运算符:它在源序列结束后,仅把「结束时刻往前推duration毫秒窗口内」产生的元素发射给下游,同时用指定的调度器(Scheduler)负责计时与缓冲队列的排空。本指南以 官方 API 文档 为主线,结合仓库中的核心实现(src/core/linq/observable/takelastwithtime.js)、模块化实现(src/modular/observable/takelastwithtime.js)与两套单元测试,讲解其调用签名、默认调度器行为、内部滑动缓冲队列原理、边界条件以及完整的可运行示例。读完本文,你将能熟练使用takeLastWithTime实现「只看最近一段时间产生的数据」这类时间窗口需求,并能读懂其源码级的时间判定逻辑。
一、功能概述
takeLastWithTime的官方语义是:
Returns elements within the specified duration from the end of the observable source sequence, using the specified schedulers to run timers and to drain the collected elements.
即:返回源序列末尾指定时间范围内产生的元素。它只关心「时间」,不关心「个数」——与按个数截取末尾元素的takeLast系列形成对比。源序列一旦发出onCompleted通知,运算符就会以此刻为基准,从内部队列中取出「结束时刻 − duration ≤ 元素产生时刻 ≤ 结束时刻」的所有元素,依次发射后完成。
与它互补的姊妹运算符是 skipLastWithTime:后者跳过末尾时间窗口内的元素、只发射窗口之外的旧元素,二者内部都使用「带时间戳的滑动队列」这一相同的数据结构。
二、方法签名与参数说明
Rx.Observable.prototype.takeLastWithTime(duration, [timeScheduler], [loopScheduler])官方文档给出的参数定义如下:
| 参数 | 类型 | 说明 | 默认值 |
|---|---|---|---|
duration | Number | 从序列末尾开始取元素的持续时间(毫秒),即时间窗口宽度 | 必填 |
timeScheduler | Scheduler | 负责运行定时器的调度器 | Rx.Scheduler.timeout |
loopScheduler | Scheduler | 负责排空(drain)已收集元素的调度器 | Rx.Scheduler.currentThread |
返回值:Observable—— 一个包含源序列末尾指定时长内元素的 observable 序列。
一个值得注意的文档与实现差异
需要指出的是,当前仓库的实际实现只接收一个调度器参数。核心实现位于 src/core/linq/observable/takelastwithtime.js:
observableProto.takeLastWithTime = function (duration, scheduler) { isScheduler(scheduler) || (scheduler = defaultScheduler); return new TakeLastWithTimeObservable(this, duration, scheduler); };其 JSDoc 注释同样只声明了[scheduler]一个可选参数,默认值为Rx.Scheduler.timeout。模块化版本(src/modular/observable/takelastwithtime.js)也保持一致,只是在默认值上使用Scheduler.async:
module.exports = function takeLastWithTime (source, duration, scheduler) { Scheduler.isScheduler(scheduler) || (scheduler = Scheduler.async); return new TakeLastWithTimeObservable(source, duration, scheduler); };而Rx.Scheduler.timeout与Rx.Scheduler.async实为同一调度器。从 src/core/concurrency/defaultscheduler.js 可以看到它们的绑定关系:
var defaultScheduler = Scheduler['default'] = Scheduler.async = new DefaultScheduler();因此,实践中直接使用takeLastWithTime(duration)或takeLastWithTime(duration, scheduler)即可;文档中列出的第二个loopScheduler参数在当前实现中并不生效,它以单一scheduler同时承担「取时间戳」「驱动滑动窗口」与「排空队列」的职责。
三、底层实现原理:带时间戳的滑动缓冲队列
takeLastWithTime的完整实现由两个类协作完成,全部代码位于 src/core/linq/observable/takelastwithtime.js。
3.1 外层 Observable:订阅时挂载观察者
var TakeLastWithTimeObservable = (function (__super__) { inherits(TakeLastWithTimeObservable, __super__); function TakeLastWithTimeObservable(source, d, s) { this.source = source; this._d = d; this._s = s; __super__.call(this); } TakeLastWithTimeObservable.prototype.subscribeCore = function (o) { return this.source.subscribe(new TakeLastWithTimeObserver(o, this._d, this._s)); }; // ... }(ObservableBase));它继承自ObservableBase,只做一件事:订阅发生时,把下游观察者o、时长_d和调度器_s一并封装进内部观察者,再订阅源序列。
3.2 核心观察者:next阶段维护滑动窗口
var TakeLastWithTimeObserver = (function (__super__) { inherits(TakeLastWithTimeObserver, __super__); function TakeLastWithTimeObserver(o, d, s) { this._o = o; this._d = d; this._s = s; this._q = []; __super__.call(this); } TakeLastWithTimeObserver.prototype.next = function (x) { var now = this._s.now(); this._q.push({ interval: now, value: x }); while (this._q.length > 0 && now - this._q[0].interval >= this._d) { this._q.shift(); } }; // ... }(AbstractObserver));这里体现了整个运算符最关键的设计:
- 带时间戳入队:每个到达的元素
x不是直接进队列,而是以调度器时钟this._s.now()打上时间戳,包装成{ interval: now, value: x }对象。因此时间基准完全由调度器决定——换用TestScheduler即可在虚拟时间下进行确定性测试。 - 队首过期即淘汰:每次入队后,只要队首元素的时间戳距今已超过
duration(即now - 队首.interval >= duration),就从队首弹出。这保证队列中永远只保留「最近duration毫秒内」到达的元素,是一个真正的滑动窗口,而不是等到完成时才一次性筛选。 - 队尾单调性假设:淘汰逻辑只检查队首,隐含假设元素按时间单调到达;对乱序源序列,
takeLastWithTime并不保证按真实到达时间重新排序。
3.3 完成阶段:按窗口放行元素
TakeLastWithTimeObserver.prototype.completed = function () { var now = this._s.now(); while (this._q.length > 0) { var next = this._q.shift(); if (now - next.interval <= this._d) { this._o.onNext(next.value); } } this._o.onCompleted(); };onCompleted到达时,以当前时刻为基准再次筛选:队列中所有「结束时刻 − 产生时刻 ≤ duration」的元素按入队顺序(即到达顺序)依次onNext发射,最后onCompleted。注意此处的判定是<=(包含窗口边界),而next中淘汰过期元素用的是>=,两处边界语义正好互补,避免边界元素被误删。
错误路径则直接透传:error不做任何缓冲,立即onError(e)转发给下游。
3.4 三阶段行为总览
| 源事件 | 运算符行为 | 下游观察 |
|---|---|---|
onNext(x) | 打时间戳入队,弹出队首超过duration的旧元素 | 不发射任何元素 |
onError(e) | 清空队列语义不适用,直接透传 | onError(e) |
onCompleted() | 以当前时刻为界,发射窗口内全部元素 | 依次onNext后onCompleted |
四、完整可运行示例(官方示例)
官方文档给出了如下可直接在浏览器控制台或 Node 中运行的示例,完整代码见 doc/api/core/operators/takelastwithtime.md:
var source = Rx.Observable.timer(0, 1000) .take(10) .takeLastWithTime(5000); var subscription = source.subscribe( function (x) { console.log('Next: ' + x); }, function (err) { console.log('Error: ' + err); }, function () { console.log('Completed'); }); // => Next: 5 // => Next: 6 // => Next: 7 // => Next: 8 // => Next: 9 // => Completed运行分析:
Rx.Observable.timer(0, 1000).take(10)产生元素0, 1, 2, …, 9,分别在第 0、1、2、…、9 秒发出,第 9 秒全部发完后onCompleted;takeLastWithTime(5000)要求只保留「完成时刻前 5 秒内」的元素。完成发生在第 9 秒,窗口起点为第 4 秒,因此0–4全部被淘汰,5–9全部保留并依次输出;- 输出顺序与源序列到达顺序一致(
5, 6, 7, 8, 9),随后Completed。
五、引入方式与发布产物
该运算符随 RxJS v4 的多个构建产物一起发布。官方文档列出的分发文件与引用关系如下(对应本仓库modules/目录下的实际产物):
- 全量构建:modules/rx-core/rx.core.js 之外的
rx.all.js、rx.all.compat.js; - 时间运算符集合:
rx.time.js(对应 modules/rx-lite-time/rx.lite.time.js 等同名产物); - 精简版:
rx.lite.js、rx.lite.compat.js(对应 modules/rx-lite/rx.lite.js); - 前置依赖:若单独使用
rx.time.js,必须先加载核心rx.js/rx.compat.js,因为takeLastWithTime依赖ObservableBase、AbstractObserver等基础设施; - 包管理:NPM 包为
rx;NuGet 包为RxJS-All、RxJS-Time、RxJS-Lite(对应仓库 nuget 目录中的RxJS-All.nuspec、RxJS-Time.nuspec、RxJS-Lite.nuspec)。
在浏览器中按顺序引入核心与时间模块后,即可通过Rx.Observable.prototype.takeLastWithTime调用。
六、单元测试与边界行为验证
仓库提供了两套测试:经典 QUnit 版位于 tests/observable/takelastwithtime.js,模块化 tape 版位于 src/modular/test/takelastwithtime.js,二者用例结构一致,均用TestScheduler在虚拟时间下验证行为。以下为 QUnit 版的关键用例归纳:
| 测试名 | 输入(虚拟时间) | duration | 期望结果 |
|---|---|---|---|
zero 1/zero 2 | 210/220/230 产生元素,230 完成 | 0 | 只onCompleted(230),无任何元素 |
some 1 | 210→1,220→2,230→3,240 完成 | 25 | onNext(240,2)、onNext(240,3)、onCompleted(240) |
some 2 | 210→1,220→2,230→3,300 完成 | 25 | 仅onCompleted(300)(元素全部过期) |
some 3 | 210–290 每 10ms 一个元素,300 完成 | 45 | onNext(300,6..9)、onCompleted(300) |
some 4 | 210–300 稀疏到达,350 完成 | 25 | 仅onCompleted(350)(元素全部过期) |
all | 210→1,220→2,230 完成 | 50 | onNext(230,1)、onNext(230,2)、onCompleted(230) |
error | 210 抛错 | 50 | onError(210, error),无缓冲直接透传 |
never | 永不完成 | 50 | 无任何消息,订阅持续到 1000 |
从这些用例可以提炼出四个关键边界结论:
duration = 0时结果为空:now - interval >= 0使所有已入队元素立即被淘汰(测试zero 1/2);- 元素在完成瞬间是否保留,取决于「完成时刻 − 产生时刻」是否
<= duration(测试some 1中 220 时刻的元素被保留,而 210 时刻的元素被淘汰,边界精确到虚拟时间刻度); - 窗口完全落在过去时,序列「静默完成」:一个元素都不发射,仅产生
onCompleted(测试some 2、some 4); - 错误不参与缓冲:源序列报错时立即转发,队列中已收集的元素被丢弃(测试
error)。
七、典型应用场景
结合上述行为特征,takeLastWithTime适合以下场景:
- 「只看最近 N 秒数据」的仪表盘/监控面板:订阅一个持续产生事件的热序列,事件结束后只关心收尾阶段最近几秒的指标,忽略较早的历史数据;
- 回放与审计:在流结束时,只取出「最近一段时间内发生的事件」用于审计或日志裁剪;
- 与
TestScheduler结合做确定性时间测试:由于时间戳全部取自注入的调度器,可以在虚拟时间里精确断言「哪些元素落在结束时刻前duration窗口内」; - 按时间而非个数截取末尾:当你不关心元素个数、只关心时间跨度时(例如「只保留最后 5 秒内的采样点」),它比按个数取末尾更贴合需求。
需要提醒的是:它是一个被动缓存型运算符——在源序列完成之前,它不会向下游发射任何元素,所有候选元素都会先进入内部队列,因此不适合需要实时输出的场景;如果你希望「实时发射且丢弃过早的元素」,应改用时间滑窗类运算符(如bufferWithTime、windowWithTime等,见 doc/api/core/operators 目录下的相关文档)。
八、相关资源
- 官方 API 文档:doc/api/core/operators/takelastwithtime.md
- 核心实现:src/core/linq/observable/takelastwithtime.js
- 模块化实现:src/modular/observable/takelastwithtime.js
- 单元测试(QUnit):tests/observable/takelastwithtime.js
- 单元测试(tape):src/modular/test/takelastwithtime.js
- 互补运算符:skipLastWithTime 文档
- 默认调度器定义:src/core/concurrency/defaultscheduler.js
- 后端
【免费下载链接】RxJS
The Reactive Extensions for JavaScript
相关推荐
RxJS 4 `skipLastWithTime` 操作符深度解析:按时间窗口跳过序列末尾元素
RxJS 4 skipLastWithTime 操作符深度解析:按时间窗口跳过序列末尾元素 skipLastWithTime duration, schedul
后端RxJS 4 运算符详解:takeLast(count) 从序列末尾截取元素的缓冲式实现原理与实战
RxJS 4 运算符详解:takeLast count 从序列末尾截取元素的缓冲式实现原理与实战 本指南以 RxJS v4(Reactive Extension
后端RxJS 4 `takeLastBuffer` 操作符深入解析:从序列末尾一次性提取指定数量元素的数组
RxJS 4 takeLastBuffer 操作符深入解析:从序列末尾一次性提取指定数量元素的数组 takeLastBuffer 是 RxJS v4 中一个看似
后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考