如何用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
相关产品推荐
相关产品推荐

