RxJS:如何重复Observable直至另一个Observable发出且不中断当前任务
RxJS实现重复Observable直到指定事件触发(不中断当前流)
需求明确:重复运行一个Observable,直到另一个Observable发出值,但仅停止后续重复操作,不能中断正在执行的Observable。
之前尝试用takeUntil的写法存在问题:
const stream$ = interval(1).pipe(takeUntil(timer(1000))); stream$.pipe(repeat(), takeUntil(timer(1500)));
这段代码里,takeUntil(timer(1500))会直接终止整个流,导致正在运行的stream$(原本需要运行1000ms)提前结束,无法完成预期的两次完整运行。
我们需要的是类似repeatUntil的效果:
const stream$ = interval(1).pipe(takeUntil(timer(1000))); stream$.pipe(repeatUntil(timer(1500)));
解决方案:用repeatWhen + takeUntil组合实现
RxJS原生没有repeatUntil操作符,但可以通过repeatWhen实现相同逻辑,它能控制重复时机且不会中断当前正在执行的流:
import { interval, timer, repeatWhen, takeUntil, tap } from 'rxjs'; // 定义单次运行的流:运行1000ms后完成 const stream$ = interval(1).pipe( takeUntil(timer(1000)), tap({ complete: () => console.log('单次流执行完成') }) ); // 实现重复直到timer(1500)发出值,且不中断当前流 stream$.pipe( repeatWhen(completionNotifications => completionNotifications.pipe(takeUntil(timer(1500))) ) ).subscribe({ next: value => console.log('输出值:', value), complete: () => console.log('整个重复流终止') });
逻辑说明
repeatWhen会在前一次流完成后,根据传入的通知流决定是否重复:这里的通知流是前一次流的完成信号- 在通知流上添加
takeUntil(timer(1500)),意味着当timer(1500)发出值后,通知流会终止,repeatWhen就不会再触发下一次重复 - 当前正在运行的
stream$会不受影响,直到自己的takeUntil(timer(1000))触发完成
执行流程示意图

内容的提问来源于stack exchange,提问作者Darko Mitic
相关产品推荐
相关产品推荐

