如何让Observable等待async回调并按订阅顺序执行?
我原本以为搜搜就能找到答案,结果没找到,只好来求助。
先看这段代码:
const something = async (): Promise<void> => { return new Promise<void>((resolve, _reject) => { setTimeout(resolve, 500); }); }; const s = new Subject<number>(); s.subscribe(async (n) => { await something(); console.log(n); }); s.subscribe(async (n) => { await something(); console.log(n + 1); }); s.next(1);
这段代码能保证始终输出:
1 2
这个顺序吗?实际上因为两个Promise都是500ms延迟,大概率会按这个顺序输出,但没有任何保证。如果把代码改成这样:
const something = async (): Promise<number> => { return new Promise<number>((resolve, _reject) => { const waitFor = Math.random() * 1000; setTimeout(() => resolve(waitFor), waitFor); }); }; const s = new Subject<number>(); s.subscribe(async (n) => { const delay = await something(); console.log(n + " / " + delay); }); s.subscribe(async (n) => { const delay = await something(); console.log((n + 1) + " / " + delay); }); s.next(1);
就会出现第二个订阅者有时比第一个先完成的情况,完全看各自的随机延迟时长。
我需要让订阅者严格按订阅顺序执行,但又不能用阻塞回调——因为回调必须是async的,要支持await操作。有没有办法让Observable等前一个订阅者完成后,再调用下一个?
要实现订阅者按顺序串行执行,核心是把每个订阅者的异步回调转换成Observable的串行流,而非让它们各自独立触发异步操作。这里有几种可行的方式:
1. 用concatMap合并订阅逻辑(推荐)
把原本的多个subscribe合并成一个,用concatMap依次处理每个异步任务,保证顺序:
const something = async (): Promise<number> => { const waitFor = Math.random() * 1000; return new Promise<number>((resolve) => setTimeout(() => resolve(waitFor), waitFor)); }; const s = new Subject<number>(); s.pipe( concatMap(n => // 第一个订阅者的逻辑包装成Observable from(something()).pipe( tap(delay => console.log(n + " / " + delay)) ) ), concatMap(n => // 第二个订阅者逻辑,等第一个完成后执行 from(something()).pipe( tap(delay => console.log((n + 1) + " / " + delay)) ) ) ).subscribe(); s.next(1);
不过这种方式需要把多个订阅逻辑合并到同一个管道里,若订阅者是动态添加的,灵活性会受限。
2. 维护一个串行执行队列
如果需要支持动态添加订阅者,可以自己维护一个Promise队列,让每个订阅者的异步任务加入队列,保证串行执行:
const something = async (): Promise<number> => { const waitFor = Math.random() * 1000; return new Promise<number>((resolve) => setTimeout(() => resolve(waitFor), waitFor)); }; const s = new Subject<number>(); // 维护串行执行的Promise链 let executionQueue = Promise.resolve(); // 封装订阅方法,自动将回调加入串行队列 const serialSubscribe = (callback: (n: number) => Promise<void>) => { return s.subscribe(async (n) => { executionQueue = executionQueue.then(async () => { await callback(n); }); }); }; // 按顺序添加订阅者 serialSubscribe(async (n) => { const delay = await something(); console.log(n + " / " + delay); }); serialSubscribe(async (n) => { const delay = await something(); console.log((n + 1) + " / " + delay); }); s.next(1);
这种方式下,无论何时添加订阅者,它们的回调都会按订阅顺序串行执行,前一个完成后才会触发下一个。
3. 自定义RxJS操作符
如果想更贴合RxJS的风格,可以自定义一个操作符处理串行订阅:
import { Observable, OperatorFunction } from 'rxjs'; const serialSubscribeOperator = <T>(callbacks: Array<(value: T) => Promise<void>>): OperatorFunction<T, void> => { return (source) => new Observable<void>((subscriber) => { let queue = Promise.resolve(); const subscription = source.subscribe({ next: (value) => { queue = queue.then(async () => { for (const callback of callbacks) { await callback(value); } }).catch(err => subscriber.error(err)); }, error: err => subscriber.error(err), complete: () => subscriber.complete() }); return subscription; }); }; // 使用方式 const s = new Subject<number>(); s.pipe( serialSubscribeOperator([ async (n) => { const delay = await something(); console.log(n + " / " + delay); }, async (n) => { const delay = await something(); console.log((n + 1) + " / " + delay); } ]) ).subscribe(); s.next(1);
这个操作符会把传入的回调数组按顺序串行执行,每个值都会依次触发所有回调,前一个完成再执行下一个。
RxJS的默认订阅机制是并行触发所有订阅者,每个订阅者的异步操作完全独立,所以无法保证执行顺序。要实现串行,必须把多个异步任务纳入同一个串行执行流——要么通过RxJS的管道操作符(如concatMap)把任务串联,要么通过外部Promise队列控制执行顺序,确保前一个任务完成后再启动下一个。
内容的提问来源于stack exchange,提问作者Lucio Crusca

