RxJS调用Subject的next方法时如何等待所有异步订阅执行完成
问题根因
- 你当前使用的写法无法实现预期输出,核心原因有两个:
- RxJS的
Subject.next()是同步执行所有注册的订阅回调的,默认不会消费回调返回的Promise,也不会等待回调内部的异步逻辑执行完成,执行完所有回调后会立刻向下执行后续同步代码。 setTimeout属于异步宏任务,会被推入事件队列等待当前同步代码全部执行完成后才会执行,因此console.log('<<< FINISH')必然先于两个定时器的输出。
- RxJS的
实现方案
你可以通过两种方式实现预期效果:
方案1:使用RxJS操作符管理异步流(推荐)
将异步任务封装为Observable流,通过操作符控制执行顺序,等待所有异步任务完成后再输出结束日志:
import { Subject, concat, from } from "rxjs"; import { concatMap, last } from "rxjs/operators"; const subject = new Subject(); // 配置异步任务执行流 const taskFlow$ = subject.pipe( concatMap(() => concat( // 第一个异步任务 from(new Promise(res => { setTimeout(() => { console.log('!! 1'); res(void 0); }, 500); })), // 第二个异步任务,会等待前一个任务完成后才执行 from(new Promise(res => { setTimeout(() => { console.log('!! 2'); res(void 0); }, 1000); })) )), // 仅在所有异步任务全部执行完成后触发订阅 last() ) console.log('>>> START'); // 所有任务执行完成后输出结束日志 taskFlow$.subscribe(() => { console.log('<<< FINISH'); }); subject.next();
该方案符合RxJS的响应式编程规范,后续扩展异步任务时成本极低。
方案2:手动维护异步计数器(轻量场景可选)
如果不想引入额外RxJS操作符,也可以手动计数判断所有异步任务是否执行完成:
import { Subject } from "rxjs"; const subject = new Subject(); // 异步任务总数量 const totalTask = 2; let finishedTask = 0; const tryFinish = () => { if (++finishedTask === totalTask) { console.log('<<< FINISH'); } } subject.subscribe(() => { setTimeout(() => { console.log('!! 1'); tryFinish(); }, 500); }) subject.subscribe(() => { setTimeout(() => { console.log('!! 2'); tryFinish(); }, 1000); }) console.log('>>> START') subject.next();
该方案逻辑简单但可维护性差,异步任务数量变动时需要手动修改计数器配置。
内容的提问来源于stack exchange,提问作者Evgeny Naumov
相关产品推荐
相关产品推荐

