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

如何用RxJS实现Observable队列化,解决认证系统异步并发问题?

RxJS实现Observable串行队列解决授权操作并发问题

我在授权系统中遇到异步并发问题,当前用RxJS的Observable处理相关操作。需求是登录、注册、登出操作同一时间只能处理一个,所有这类操作必须按顺序等待前一个完成后再执行。

问题出在注册流程:注册函数需要先在凭证库注册账户信息,再把账户数据写入账户数据库,这是两个依次执行的请求。但认证库会在凭证确认后触发一个事件,这个事件可能在账户数据写入数据库前就被触发,导致事件里的读取操作拿不到数据。我无法修改事件触发逻辑,所以需要一个队列机制,让每个请求排队等待前一个完成后再执行,同时用标记标识数据正在写入,让读取操作等写入完成后再执行。

伪代码如下:

// 这个函数会在signUpWithEmail完成后、writeDbData执行前被调用
function onSignInEvent(userId): void {
    readDbData(userId).subscribe(); // 此时注册还没完成,读取不到数据
}

function signUp(email, password): Observable<Account> {
    return from(signUpWithEmail(email, password)).pipe(
        // 这里凭证已注册,事件会被触发
        map(userId => createAccountData(userId)),
        switchMap(account => writeDbData(account)) // 这里才会把数据写入数据库
    );
}

实现思路

用RxJS的Subject作为操作队列的调度中心,配合concatMap操作符实现串行执行——concatMap会严格等待前一个Observable完成后,再订阅执行下一个Observable,天然符合队列要求。所有需要串行的操作(注册、登录、登出,以及事件触发的读取操作)都通过这个队列调度,确保不会出现并发冲突。

具体实现代码

import { Subject, Observable, from, concatMap, tap, catchError } from 'rxjs';

// 创建一个Subject作为操作队列,每个待执行的操作都是一个返回Observable的函数
const authOperationQueue$ = new Subject<() => Observable<any>>();

// 订阅队列,用concatMap保证串行执行
authOperationQueue$.pipe(
    concatMap(operation => operation()),
    // 可选:处理错误,避免单个操作失败导致整个队列终止
    catchError((err, caught) => {
        console.error('授权操作失败:', err);
        return caught;
    })
).subscribe();

// 修改注册函数,将整个注册流程加入队列
function signUp(email, password): Observable<Account> {
    // 创建Promise传递注册结果,适配队列的异步调度逻辑
    const resultPromise = new Promise<Account>((resolve, reject) => {
        const operation = () => from(signUpWithEmail(email, password)).pipe(
            map(userId => createAccountData(userId)),
            switchMap(account => writeDbData(account)),
            tap({
                next: resolve,
                error: reject
            })
        );
        authOperationQueue$.next(operation);
    });
    return from(resultPromise);
}

// 修改onSignInEvent,将读取操作加入队列
function onSignInEvent(userId): void {
    const operation = () => readDbData(userId).pipe(
        tap(data => {
            // 此时注册的写操作已完成,能读取到正确数据
            console.log('读取到账户数据:', data);
        })
    );
    authOperationQueue$.next(operation);
}

关键说明

  • 队列调度:所有需要串行的操作都包装成返回Observable的函数,通过authOperationQueue$.next()加入队列,concatMap会按顺序执行每个操作,前一个完成后才会执行下一个。
  • 注册结果传递:用Promise适配队列的异步调度,保证调用signUp的地方能正常订阅到最终的注册结果。
  • 错误容错:加入catchError并返回caught,确保单个操作失败不会导致整个队列停止,后续操作依然能正常执行。
  • 读取操作等待:onSignInEvent里的读取操作加入队列后,会自动等待之前的注册写操作完成,此时数据库已写入数据,能读取到正确结果。

内容的提问来源于stack exchange,提问作者magg

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 21:05:02