RxJS 4 的 take 操作符详解:从序列头部精确截取 N 个元素
RxJS 4 的 take 操作符详解:从序列头部精确截取 N 个元素
Rx.Observable.prototype.take(count, [scheduler]) 是 RxJS(The Reactive Extensions for JavaScript,v4 分支)中用于从可观察序列头部截取指定数量连续元素的核心操作符。它常用于限制数据流长度、实现"只取前 N 个结果"、或与 interval/range 等无限序列配合构造有界场景。读完本文,你将掌握 take 的完整签名与参数语义、count 为 0/负数时的边界行为、底层 TakeObservable 的实现原理,以及通过仓库自带测试用例验证的完整时序语义。
本文基于当前仓库 doc/api/core/operators/take.md 展开,并对照源码 src/core/perf/operators/take.js、模块化实现 src/modular/observable/take.js 与单元测试 tests/observable/take.js 进行纵深剖析。
一、API 签名与参数说明
Rx.Observable.prototype.take(count, [scheduler])
功能:返回从源可观察序列起始位置开始的指定数量的连续元素(Returns a specified number of contiguous elements from the start of an observable sequence)。
参数(Arguments)
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
count |
Number |
是 | 要返回的元素数量(The number of elements to return)。 |
scheduler |
Scheduler |
否 | 当 count 被设置为 0 时,用于产生 onCompleted 消息的调度器(Scheduler used to produce an onCompleted message in case count is set to 0)。 |
返回值(Returns)
(Observable):一个包含输入序列中指定数量元素的可观察序列。文档原文描述为"包含输入序列中指定索引之前(含该索引)元素的可观察序列",在实际语义上即源序列的前 count 个元素。
前置条件
- 无(Prerequisites: None),
take是纯基础操作符,不依赖其他模块。
二、快速上手示例
文档给出的官方示例使用 Rx.Observable.range(0, 5).take(3),完整可运行代码如下:
var source = Rx.Observable.range(0, 5)
.take(3);
var subscription = source.subscribe(
function (x) {
console.log('Next: ' + x);
},
function (err) {
console.log('Error: ' + err);
},
function () {
console.log('Completed');
});
// => Next: 0
// => Next: 1
// => Next: 2
// => Completed
执行结果清晰地展示了两个关键行为:
range(0, 5)本应依次发射0, 1, 2, 3, 4共 5 个元素;- 经过
.take(3)后,只接收到前 3 个元素0, 1, 2,并且在第 3 个元素之后立即触发Completed——即take会主动终止订阅,而不是等源序列自然结束。
三、边界行为:count 为 0 与负数的处理
take 对 count 的边界值有明确的规范化处理,这在源码中体现得淋漓尽致。
3.1 count 为负数:抛出异常
在 src/core/perf/operators/take.js 的入口实现中:
observableProto.take = function (count, scheduler) {
if (count < 0) { throw new ArgumentOutOfRangeError(); }
if (count === 0) { return observableEmpty(scheduler); }
return new TakeObservable(this, count);
};
当 count < 0 时,直接抛出 ArgumentOutOfRangeError(参数越界错误),拒绝负数的非法取值。模块化版本 src/modular/observable/take.js 行为完全一致,只是错误类来自内部模块 errors.ArgumentOutOfRangeError:
module.exports = function (source, count, scheduler) {
if (count < 0) { throw new errors.ArgumentOutOfRangeError(); }
if (count === 0) { return empty(scheduler); }
return new TakeObservable(source, count);
};
3.2 count 为 0:立即完成
当 count === 0 时,take 不会创建任何订阅,而是直接返回一个"空序列" observableEmpty(scheduler)。这正是第二个可选参数 scheduler 唯一发挥作用的地方——用指定的调度器异步(或按调度策略)发送 onCompleted 消息。默认情况下,observableEmpty 使用 Scheduler.immediate 立即完成;传入 Rx.Scheduler.async 或 Rx.Scheduler.timeout 则可让完成通知延后到调度队列中。
类型定义文件 ts/core/linq/observable/take.ts 也给出了对应的 TS 重载用法示例:
take(count: number, scheduler?: IScheduler): Observable<T>;
// 使用示例
var res = source.take(5);
var res = source.take(0, Rx.Scheduler.timeout);
四、底层实现原理:TakeObservable 与 TakeObserver
take 的正常路径(count > 0)通过一对内部类完成:TakeObservable(序列工厂)与 TakeObserver(订阅观察者)。
4.1 TakeObservable:订阅即包装
var TakeObservable = (function(__super__) {
inherits(TakeObservable, __super__);
function TakeObservable(source, count) {
this.source = source;
this._count = count;
__super__.call(this);
}
TakeObservable.prototype.subscribeCore = function (o) {
return this.source.subscribe(new TakeObserver(o, this._count));
};
...
}(ObservableBase));
TakeObservable 继承自 ObservableBase,保存源序列 source 与目标数量 _count。当外部订阅发生时(subscribeCore),它把下游观察者 o 连同 _count 一起包装成 TakeObserver,并订阅到源序列上。整个截取逻辑完全封装在观察者内部。
4.2 TakeObserver:剩余计数器驱动终止
function TakeObserver(o, c) {
this._o = o;
this._c = c;
this._r = c; // 剩余可接收数量
AbstractObserver.call(this);
}
TakeObserver.prototype.next = function (x) {
if (this._r-- > 0) {
this._o.onNext(x);
this._r <= 0 && this._o.onCompleted(); // 数量耗尽立即完成
}
};
TakeObserver.prototype.error = function (e) { this._o.onError(e); };
TakeObserver.prototype.completed = function () { this._o.onCompleted(); };
这里的关键设计有三点:
- 先判断后转发:
if (this._r-- > 0)使用后置自减,每收到一个元素,剩余配额_r减 1;只有配额未耗尽时才向下游转发onNext(x)。 - 配额耗尽立即完成:
this._r <= 0 && this._o.onCompleted()意味着在转发恰好第count个元素之后,立刻向下游发送onCompleted,从而主动切断与源序列的连接。这就是示例中"输出 3 个元素后马上 Completed"的底层原因。 - 错误与提前完成透传:
error与completed原样转发——若源序列在配额耗尽前出错,错误照常传播;若源序列比count短而先行完成,则自然完成,不存在"补足数量"或"等待"逻辑。
4.3 与 AbstractObserver 的关系
TakeObserver 继承自 AbstractObserver(基类位于 src/core/abstractobserver.js),其职责是统一实现 onNext/onError/onCompleted 的调用规约(例如完成/错误后屏蔽后续通知等防御逻辑),而 take 只需覆写 next/error/completed 三个回调即可专注业务规则。这也是 RxJS 4 中大量操作符共用的观察者模式。
五、时序语义:来自测试用例的验证
仓库为 take 提供了非常完备的虚拟时间测试套件 tests/observable/take.js,使用 Rx.TestScheduler 精确断言消息时序与订阅生命周期。这些用例是理解 take 行为的最佳"活文档"。
5.1 三种完成时机
测试构造了一个热序列(Hot Observable),在 210 到 630 之间连续发射 17 个元素,并分别用不同 count 验证三种场景:
| 场景 | 调用 | 结果 | 依据 |
|---|---|---|---|
| 完成在源之后(complete after) | xs.take(20) |
全部 17 个元素照单全收,跟随源在 690 完成 |
tests/observable/take.js |
| 完成与源同步(complete same) | xs.take(17) |
恰好收到第 17 个元素(630)时立即完成,不等源在 690 完成 |
tests/observable/take.js |
| 完成在源之前(complete before) | xs.take(10) |
收到第 10 个元素(415)即完成,订阅在 415 提前解除 |
tests/observable/take.js |
特别值得注意 take(17) 用例:源序列发射 17 个元素后,take 在第 17 个元素处立即 onCompleted,因此订阅区间为 (200, 630),早于源序列自身的完成时刻 690——这说明 take 的完成信号由自身配额驱动,与源序列的终止时刻解耦。
5.2 错误传播语义
测试还覆盖了错误场景(error after / error same / error before):
- 当
count大于源序列元素总数时(如take(20)),源在690抛出的onError会被原样透传,下游观察到错误; - 当
count恰好等于元素总数时(如take(17)),take在630自行完成,源序列后续的错误(690)不会再传播到下游——订阅早已在630解除,这印证了"提前完成即彻底隔离"的特性; - 当
count小于元素总数时(如take(3)),第三个元素后即完成(270),源在690的错误同样被隔离。
5.3 中途取消订阅(Dispose)
take(3)且在250时刻释放订阅:只收到210、230两个元素,订阅区间为(200, 250);- 释放时刻
400晚于完成时刻270:三个元素全部收到并于270完成,订阅区间仍为(200, 270)——完成后订阅即终止,不存在悬挂。
六、典型组合用法
take 常与以下操作符搭配,构成实用模式(以下用法均可直接运行):
限制无限序列的长度:
// 每秒发射一个递增整数,只取前 5 个
var source = Rx.Observable.interval(1000)
.take(5)
.subscribe(
function (x) { console.log('Next: ' + x); },
function (e) { console.log('Error: ' + e); },
function () { console.log('Completed'); });
// => Next: 0, 1, 2, 3, 4, Completed(约 5 秒后完成)
与 takeWhile / takeUntil 对照: 同族操作符中,takeWhile 按谓词条件截取、takeUntil 按另一个序列的信号截取,而 take 是纯粹的按数量截取;反向的"跳过前 N 个"则由 skip 系列操作符承担。选择依据是:数量确定用 take,条件确定用 takeWhile,事件驱动用 takeUntil。
七、发布产物与获取方式
take 随多个发布产物分发,无需额外引入(Prerequisites: None):
- 源码:本仓库的实现位于 src/core/perf/operators/take.js(文档标注的
src/core/linq/observable/take.js在当前仓库中已由 perf 目录下的实现取代),模块化独立版本见 src/modular/observable/take.js; - Dist 构建:
rx.js、rx.compat.js、rx.all.js、rx.all.compat.js、rx.lite.js、rx.lite.compat.js等(本仓库对应产物位于 modules 目录,如 modules/rx-lite/rx.lite.js); - NPM 包:
rx; - NuGet 包:
RxJS-All、RxJS-Main、RxJS-Lite(对应清单见 nuget 目录下的.nuspec文件); - 单元测试:tests/observable/take.js,可与
QUnit+Rx.TestScheduler一起运行验证; - 类型定义:ts/core/linq/observable/take.ts。
八、总结
Rx.Observable.prototype.take(count, [scheduler]) 以极小的实现成本(一个计数器 _r)提供了明确、高效的"取前 N 个元素"语义,其行为可概括为:
count < 0→ 抛出ArgumentOutOfRangeError;count === 0→ 返回observableEmpty(scheduler),用可选调度器发送完成通知;count > 0→ 转发前count个元素,配额耗尽即刻onCompleted并解除订阅;- 错误在配额耗尽前到达则透传,配额耗尽后到达则被隔离。
理解 take 的提前终止与订阅解除语义,是正确使用 RxJS 处理有限/无限序列、避免资源泄漏的基础,也是读懂 takeWhile、takeUntil 等同族操作符的起点。