RxJS如何实现待指定Observable完成后发送源流剩余值
实现方案
要获取stream$中0之后的所有剩余值(排除3、2、1、0),有两种可直接运行的实现方式,都能兼容你现有的ready$逻辑。
写法1:极简实现(用skipWhile)
skipWhile的逻辑和你用的takeWhile刚好相反:会持续跳过符合条件的值,直到条件返回false后,就正常发送后续所有值。由于触发停止跳过的第一个值是0,额外加skip(1)把0本身排除即可。
import { of } from 'rxjs'; import { takeWhile, skipWhile, skip } from 'rxjs/operators'; const stream$ = of(3, 2, 1, 0, -1, -2, -3); const ready$ = stream$.pipe(takeWhile((data) => data !== 0)); // 拿0之后的所有值 const afterReady$ = stream$.pipe( skipWhile(data => data !== 0), skip(1) ); // 测试 ready$.subscribe(console.log); // 输出3、2、1 afterReady$.subscribe(console.log); // 输出-1、-2、-3
写法2:和ready$生命周期绑定
如果要严格对齐「ready$完成后才发送后续值」的要求,不重复写0的判定逻辑,可以用skipUntil监听ready$的完成事件作为放行信号:
import { of, endWith, ignoreElements } from 'rxjs'; import { takeWhile, skipUntil, skip } from 'rxjs/operators'; const stream$ = of(3, 2, 1, 0, -1, -2, -3); const ready$ = stream$.pipe(takeWhile((data) => data !== 0)); const afterReady$ = stream$.pipe( skipUntil(ready$.pipe(ignoreElements(), endWith(true))), skip(1) );
边界场景适配
如果业务里0可能重复出现,需要保证只在第一次遇到0之后放行,用带状态标记的写法更稳妥,不会被后续的0干扰逻辑:
import { of } from 'rxjs'; import { scan, filter, map } from 'rxjs/operators'; const stream$ = of(3, 2, 1, 0, -1, -2, 0, -3); const afterReady$ = stream$.pipe( scan((acc, cur) => { if (cur === 0) acc.triggered = true; return { triggered: acc.triggered, value: cur }; }, { triggered: false, value: null }), filter(item => item.triggered && item.value !== 0), map(item => item.value) ); // 输出-1、-2、-3
注:
of创建的是冷数据流,多次订阅会独立执行完整序列,不存在多订阅值冲突的问题。
内容的提问来源于stack exchange,提问作者Ultra
相关产品推荐
相关产品推荐

