RxJS window操作符连续相同元素分组求和异常修复
问题根因
分组错位的核心原因有两个:
window操作符的切窗逻辑是收到边界通知时立刻关闭当前窗口、打开新窗口,但原代码中checkChange的切窗通知触发时机不对:第一个异值(第一个2)发射后,pairwise才会检测到[1,2]的差异发出切窗信号,默认订阅顺序下这个异值会先落入旧窗口,再触发切窗动作- 缺少初始开窗信号,第一个窗口的生命周期完全绑定第一次切窗信号,进一步放大了时序错位问题,最终分组变成
[[1,1,1,2],[2,1],[1]],求和结果自然不符合预期。
修复方案
提供两种可直接运行的实现,优先选择第一种,无操作符时序风险,逻辑更可控:
方案1:scan状态聚合实现(推荐)
完全自主控制分组累加逻辑,不依赖window的切窗时序,代码稳定性和可读性更高:
import { from } from 'rxjs'; import { scan, concatMap } from 'rxjs/operators'; from([1, 1, 1, 2, 2, 1, 1]) .pipe( scan((state, currNum) => { // 初始化第一个分组 if (!state) { return { currentGroupValue: currNum, groupSum: currNum, finishedSums: [] }; } // 与当前分组值相同,直接累加 if (currNum === state.currentGroupValue) { state.groupSum += currNum; return state; } // 值发生变化,将已完成的分组和存入输出队列,重置新分组 state.finishedSums.push(state.groupSum); state.currentGroupValue = currNum; state.groupSum = currNum; return state; }, null), // 追加最后一个未入队的分组和,输出最终结果 concatMap(finalState => [...finalState.finishedSums, finalState.groupSum]) ) .subscribe(console.log); // 输出:3 4 2
方案2:修复原有window切窗逻辑
如果要保留window的写法,需要补全初始开窗信号,利用connect共享源的订阅顺序保证切窗动作先于异值入窗:
import { from, merge, of } from 'rxjs'; import { connect, pairwise, filter, map, window, mergeMap, reduce } from 'rxjs/operators'; from([1, 1, 1, 2, 2, 1, 1]) .pipe( connect((numbers$) => { const windowBoundary$ = merge( of('init'), // 订阅时立刻触发第一个窗口开启 numbers$.pipe( pairwise(), filter(([prev, curr]) => prev !== curr), map(() => 'change') ) ); return numbers$.pipe( window(windowBoundary$), mergeMap(group$ => group$.pipe(reduce((acc, n) => acc + n, 0))) ); }) ) .subscribe(console.log); // 输出:3 4 2
注意:如果源是异步发射值,需要给切窗信号加
observeOn(asyncScheduler)保证切窗动作优先级高于值入窗,同步源场景下上述写法可直接运行。
内容的提问来源于stack exchange,提问作者Max Smirnov
相关产品推荐
相关产品推荐

