为何RXJS中外层takeUntil无法终止switchMap内的timer?
takeUntil来终止timer? 组件代码
private takeControlTaskRunning$ = new BehaviorSubject<boolean>(false); private initTakeControlTask() { this.takeControlTaskRunning$ .pipe( takeUntil(this.dependencies.destroyed$), switchMap((isRunning) => { if (!isRunning) { return of(false); } return timer(0, this.updatePeriodMilliseconds).pipe( takeUntil(this.dependencies.destroyed$), // <--- 为何需要这行? map(() => { return true; }) ); }), filter((isRunning) => isRunning === true), ) .subscribe(async () => { await this.syncVideo(); }); } public takeControl() { this.takeControlTaskRunning$.next(true); } public stopControl() { this.takeControlTaskRunning$.next(false); }
测试代码
describe('takeControl', () => { it('runs the update task periodically', fakeAsync(async () => { // Given an controller instance const { controller, dependencies } = await create(); // When takeControl is called await controller.takeControl(); // And some time passes jasmine.clock().tick(1); // Then the update task should be running expect( dependencies.updateSync as jasmine.Spy ).toHaveBeenCalled(); destroyed$.next(); flush(); })); });
错误信息
Error: 1 periodic timer(s) still in the queue.
疑问
添加内层takeUntil可解决问题,但我不理解其必要性。按我的理解,外层takeUntil应完成整个Observable,switchMap也应终止内部timer,为何实际并非如此?
核心原因在于fakeAsync环境下的定时器管理机制,以及RxJS中订阅终止的传递逻辑:
外层takeUntil的作用边界:
外层的takeUntil(this.dependencies.destroyed$)确实会让整个pipe的订阅完成,但它的信号是向上传递的——终止的是takeControlTaskRunning$.pipe(...)这个整体的订阅,并不会直接主动终止switchMap内部的timerObservable订阅。switchMap的终止条件:
switchMap只会在源Observable(也就是takeControlTaskRunning$)发出新值时,才会取消前一个内部Observable的订阅。而你在测试中是直接调用destroyed$.next(),并没有触发takeControlTaskRunning$发出新值,所以switchMap不会主动去终止当前的timer订阅。fakeAsync的严格定时器检查:
fakeAsync会模拟所有定时器,并且要求测试结束时所有定时器必须被清理。如果内部timer没有被主动终止,它的周期性任务会一直留在fakeAsync的定时器队列中,即使外层订阅已经结束,队列里的任务依然存在,最终导致测试报错。
添加内层takeUntil(this.dependencies.destroyed$)后,当destroyed$发出信号时,会直接终止内部的timer Observable,清理掉它在fakeAsync队列中的所有周期性任务,从而解决报错问题。
内容的提问来源于stack exchange,提问作者Blasco

