You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 16:21:11