尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RxJS v4 windowWithCount 操作符详解:按元素数量将可观测序列切分为多个窗口

RxJS v4 windowWithCount 操作符详解:按元素数量将可观测序列切分为多个窗口 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])参数类型必填说明countNumber是每个窗口的长度包含的元素个数。必须大于 0否则抛出ArgumentOutOfRangeError。skipNumber否相邻两个窗口起始位置之间跳过的元素个数。如果不提供默认等于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的检查拦截抛出ArgumentOutOfRangeErrorskip缺省推导skip nullundefined或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于是产生窗口 11, 2窗口 23, 4窗口 35, 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 个元素于是产生高度重叠的滑动窗口窗口 11, 2窗口 22, 3窗口 33, 4窗口 44, 5窗口 55, 6窗口 66收尾时不完整窗口末尾空窗口这种滑动窗口模式非常适合用于移动平均、滑动统计、去重窗口等需要上下文重叠的流式计算场景。四、返回值说明窗口是 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 个、第 countskip 个、第 count2*skip 个……元素到达时恰好有一个窗口被填满而关闭开启新窗口每次计数自增后若n % skip 0则创建下一个窗口。因此第一个窗口是在第skip个元素到达后开启的。结合 3.2 的windowWithCount(2, 1)第 2 个元素到达时n2c 2-21 11 % 1 0窗口 1 关闭同时2 % 1 0开启窗口 2从而形成滑动窗口。以 3.1 的windowWithCount(2)为例第 2 个元素到达时c11 % 2 ! 0窗口 1 不关闭但2 % 2 0会开启窗口 2第 4 个元素到达时c33 % 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.jsobservableProto.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 热可观测序列覆盖三个场景basicwindowWithCount(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。模块化版windowcount.js使用tapereactiveAssert复现了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窗口长度必填0skip窗口间距缺省 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),仅供参考
返回列表