RxJS v4 windowWithCount 操作符详解:按元素数量将可观测序列切分为多个窗口
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
本文围绕 RxJS v4 的windowWithCount(别名windowCount)操作符展开,讲解如何基于元素计数信息把一个可观测序列"项目化"为零个或多个窗口(window),每个窗口本身又是一个可观测序列。你将掌握count与skip两个参数的确切语义、窗口开启与关闭的判定规则,以及它与bufferWithCount的等价关系,并能结合源码与单元测试理解其底层实现原理。
一、操作符概述
windowWithCount是 RxJS v4 中用于按元素个数切分数据流的核心操作符之一。它的工作方式与buffer系列类似,但关键区别在于:
bufferWithCount产出的是数组(缓冲区);windowWithCount产出的是可观测序列(窗口),每个窗口包含一段连续的元素,窗口之间可以重叠也可以不重叠。
该操作符在 核心实现文件 中定义,官方 API 文档位于 windowwithcount.md,完整的行为契约由 单元测试 固定。
二、API 签名与参数语义
Rx.Observable.prototype.windowWithCount(count, [skip])| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
count | Number | 是 | 每个窗口的长度(包含的元素个数)。必须大于 0,否则抛出ArgumentOutOfRangeError。 |
skip | Number | 否 | 相邻两个窗口起始位置之间跳过的元素个数。如果不提供,默认等于count,即窗口不重叠。 |
返回:一个Observable,其每个推送值都是一个窗口(也是Observable)。
从 源码第 7-16 行 可以看到完整的参数校验逻辑:
observableProto.windowWithCount = observableProto.windowCount = function (count, skip) { var source = this; +count || (count = 0); Math.abs(count) === Infinity && (count = 0); if (count <= 0) { throw new ArgumentOutOfRangeError(); } skip == null && (skip = count); +skip || (skip = 0); Math.abs(skip) === Infinity && (skip = 0); if (skip <= 0) { throw new ArgumentOutOfRangeError(); } // ... };这里有几个值得注意的实现细节:
- 一元
+强制数值转换:count会被+count先转为数值,NaN或0会被置为0; - 无穷大兜底:
Math.abs(count) === Infinity时同样置为0,随后被count <= 0的检查拦截,抛出ArgumentOutOfRangeError; skip缺省推导:skip == null(undefined或null)时回退为count,因此不传skip等价于"不重叠窗口";- 两者都必须为正整数:
count、skip小于等于 0 都会抛出ArgumentOutOfRangeError。
同时,从代码可以看出windowWithCount与windowCount是同一个函数(别名关系),使用哪一个名字效果完全相同。
三、两种调用方式与完整示例
官方文档 示例 给出了两种典型用法,下面完整复现并逐行解读。
3.1 不传 skip:窗口不重叠
/* Without a skip */ var source = Rx.Observable.range(1, 6) .windowWithCount(2) .selectMany(function (x) { return x.toArray(); }); var subscription = source.subscribe( function (x) { console.log('Next: ' + x.toString()); }, function (err) { console.log('Error: ' + err); }, function () { console.log('Completed'); }); // => Next: 1,2 // => Next: 3,4 // => Next: 5,6 // => Next: // => Completedrange(1, 6)依次发出1, 2, 3, 4, 5, 6。windowWithCount(2)表示每个窗口长度 2、窗口间距也为 2,于是产生:
- 窗口 1:
1, 2 - 窗口 2:
3, 4 - 窗口 3:
5, 6
最后还会额外推送一个空窗口(输出中单独一行的Next:),这是因为当源序列完成时,最后一个尚未填满的窗口会以"空"的形式补发并立即完成——这与bufferWithCount会过滤空缓冲区的行为形成对比(详见第六节)。
selectMany(即flatMap)在这里的作用是把每个窗口(Observable)拍平为数组后合并到主序列,便于直接打印查看窗口内容。
3.2 指定 skip:窗口重叠
/* Using a skip */ var source = Rx.Observable.range(1, 6) .windowWithCount(2, 1) .selectMany(function (x) { return x.toArray(); }); var subscription = source.subscribe( function (x) { console.log('Next: ' + x.toString()); }, function (err) { console.log('Error: ' + err); }, function () { console.log('Completed'); }); // => Next: 1,2 // => Next: 2,3 // => Next: 3,4 // => Next: 4,5 // => Next: 5,6 // => Next: 6 // => Next: // => CompletedwindowWithCount(2, 1)表示每个窗口长度 2、每次只前进 1 个元素,于是产生高度重叠的滑动窗口:
- 窗口 1:
1, 2 - 窗口 2:
2, 3 - 窗口 3:
3, 4 - 窗口 4:
4, 5 - 窗口 5:
5, 6 - 窗口 6:
6(收尾时不完整窗口) - 末尾空窗口
这种"滑动窗口"模式非常适合用于移动平均、滑动统计、去重窗口等需要上下文重叠的流式计算场景。
四、返回值说明:窗口是 Observable 而非数组
windowWithCount的返回值是一个Observable,其中每个元素都是一个窗口(Observable)。这一点与buffer系列"返回元素为数组"的语义有本质区别,意味着:
- 每个窗口可以独立订阅、独立变换(如
map、filter); - 窗口之间可以并行消费,也可以像示例中那样通过
selectMany/mergeAll合并回单一数据流; - 因为窗口是惰性推送的,下游可以通过
addRef机制延迟订阅窗口,避免在窗口关闭前就被迫消费。
在 源码第 23-29 行 中,每个窗口由Subject充当:
function createWindow () { var s = new Subject(); q.push(s); observer.onNext(addRef(s, refCountDisposable)); }窗口被创建后立即推送给下游观察者;addRef(s, refCountDisposable)将窗口与一个RefCountDisposable关联,保证当所有窗口的订阅都释放后,底层源订阅会被自动清理。
五、源码级原理剖析
windowWithCount的核心逻辑位于 windowwithcount.js。其实现依赖AnonymousObservable、SingleAssignmentDisposable、RefCountDisposable与Subject的组合:
return new AnonymousObservable(function (observer) { var m = new SingleAssignmentDisposable(), refCountDisposable = new RefCountDisposable(m), n = 0, q = []; function createWindow () { var s = new Subject(); q.push(s); observer.onNext(addRef(s, refCountDisposable)); } createWindow(); m.setDisposable(source.subscribe( function (x) { for (var i = 0, len = q.length; i < len; i++) { q[i].onNext(x); } var c = n - count + 1; c >= 0 && c % skip === 0 && q.shift().onCompleted(); ++n % skip === 0 && createWindow(); }, function (e) { while (q.length > 0) { q.shift().onError(e); } observer.onError(e); }, function () { while (q.length > 0) { q.shift().onCompleted(); } observer.onCompleted(); } )); return refCountDisposable; }, source);5.1 窗口的开启与关闭时机
- 订阅即开窗:
createWindow()在订阅开始时被立即调用一次,保证第一个窗口始终存在; - 广播:每当源发出元素
x,当前所有打开中的窗口(q数组)都会收到onNext(x); - 关闭窗口:维护计数器
n(已接收元素总数)。当n - count + 1 >= 0且(n - count + 1) % skip === 0时,队首窗口调用onCompleted()并出队。直觉上就是"第 count 个、第 count+skip 个、第 count+2*skip 个……元素到达时,恰好有一个窗口被填满而关闭"; - 开启新窗口:每次计数自增后,若
++n % skip === 0,则创建下一个窗口。因此第一个窗口是在第skip个元素到达后开启的。
结合 3.2 的windowWithCount(2, 1):第 2 个元素到达时(n=2,c = 2-2+1 = 1,1 % 1 === 0)窗口 1 关闭,同时2 % 1 === 0开启窗口 2,从而形成滑动窗口。以 3.1 的windowWithCount(2)为例:第 2 个元素到达时(c=1,1 % 2 !== 0)窗口 1 不关闭,但2 % 2 === 0会开启窗口 2;第 4 个元素到达时(c=3,3 % 2 !== 0)……以此类推,可见关闭与开启的判定是相互独立的两个条件。
5.2 错误与完成的传播语义
- 错误:源发出错误时,所有尚未关闭的窗口依次
onError,随后主观察者也收到onError; - 完成:源正常完成时,所有未关闭的窗口(包括那些未填满的)依次
onCompleted,主观察者再收到onCompleted。这正是示例输出末尾出现"空窗口/不完整窗口"的原因。
5.3 模块化版本的结构
仓库还提供了一套独立可发布的模块化实现(CommonJS 风格),对应文件为 windowcount.js。它把相同逻辑拆分为WindowCountObservable(继承ObservableBase,实现subscribeCore)与WindowCountObserver(继承AbstractObserver),参数校验(第 71-81 行)与核心版完全一致,模块入口统一由 index.js 导出。
六、与 bufferWithCount 的关系
windowWithCount并非孤立存在,它是bufferWithCount的底层基石。查看 bufferwithcount.js:
observableProto.bufferWithCount = observableProto.bufferCount = function (count, skip) { typeof skip !== 'number' && (skip = count); return this.windowWithCount(count, skip) .flatMap(toArray) .filter(notEmpty); };bufferWithCount本质上就是对windowWithCount的二次加工:
- 用
windowWithCount(count, skip)切出窗口; - 通过
flatMap(toArray)把每个窗口收拢成数组; - 通过
filter(notEmpty)过滤掉空窗口。
这也解释了前面观察到的差异:bufferWithCount不会输出空数组,而windowWithCount会推送空窗口。模块化版本 buffercount.js 的实现思路完全一致。因此,理解了windowWithCount,就同时理解了bufferWithCount的一半实现。
七、单元测试验证
行为契约由两类测试锁定:
- 经典版:tests/observable/windowwithcount.js,基于
TestScheduler+ 热可观测序列,覆盖三个场景:- basic:
windowWithCount(3, 2)下窗口为[2,3,4]、[4,5,6]、[6,7,8]、[8,9],消息断言精确到时间点(如onNext(280, '0 4')与onNext(280, '1 4')同刻发生,证明窗口重叠且并行推送),同时验证底层订阅区间为[200, 600]; - disposed:在
370时刻主动取消订阅后,只产生到'2 6'为止的消息,订阅区间变为[200, 370],验证资源释放路径; - error:源在
600时刻抛出错误,所有打开中的窗口与主观察者依次收到onError。
- basic:
- 模块化版:windowcount.js,使用
tape+reactiveAssert复现了windowCount(3, 2)的 basic 场景,保证两套实现行为一致。
这些测试是理解"何时关窗、何时开窗、窗口如何并行广播"最直观的可执行证据。
八、可用性与获取方式
根据 官方文档 的说明:
- 发布版本:
windowWithCount被编译进rx.js、rx.compat.js以及轻量扩展版rx.lite.extras.js(对应构建产物可查看 dist 目录 与 rx-lite-extras 模块); - 无前置依赖:使用该操作符不需要额外引入其他模块;
- 模块化发布:可单独引入 src/modular/observable/windowcount.js 并挂载到
Observable.prototype(参见 模块化测试 的挂载方式)。
九、小结
| 维度 | 结论 |
|---|---|
| 功能 | 按元素计数把源序列切分为零个或多个窗口(Observable) |
| 参数 | count(窗口长度,必填,>0);skip(窗口间距,缺省 = count) |
| 返回值 | Observable,每个元素是独立窗口 |
| 别名 | windowCount |
| 窗口粒度 | 可重叠(skip < count)、恰好相接(skip = count)、可留空隙(skip > count) |
| 与 buffer 关系 | bufferWithCount基于本操作符 +flatMap(toArray)+filter(notEmpty)实现 |
| 收尾行为 | 源完成时未填满的窗口与一个空窗口都会被推送并完成 |
当需要以"可观测序列"而非"数组"为粒度对数据流做滑动切分、窗口内再变换或窗口级背压控制时,windowWithCount就是 RxJS v4 中直接对应的标准答案;其参数校验、窗口生命周期与错误传播细节,均可在 实现 与 测试 中进一步验证。
【免费下载链接】RxJSThe Reactive Extensions for JavaScript项目地址: https://gitcode.com/gh_mirrors/rxj/RxJS
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考