如何实现await Subject.next 等待RxJS Subject订阅侧异步逻辑执行完成
原实现问题分析
你提供的代码无法生效的核心原因是:tap 操作符会在值下发给订阅者的瞬间触发,不会等待订阅者内部的异步逻辑执行完成,因此 wrapper.awaiter 会在 next 调用后立刻 resolved,此时订阅中的延迟逻辑还未执行,自然拿不到预期的 testValue。
可行实现方案
方案1:调整流结构,将异步逻辑移入管道(推荐,符合RxJS设计理念)
不需要自定义Subject,利用RxJS原生的高阶操作符处理异步等待逻辑,代码更易维护:
import { Subject, concatMap, from, firstValueFrom, map, Observable } from 'rxjs' // 封装触发工具 function createAwaitableTrigger<T>(handler: (input: T) => Promise<void> | Observable<void>) { const subject = new Subject<T>() // 用concatMap保证按顺序处理输入,等待前一个异步任务完成再处理下一个 const process$ = subject.pipe( concatMap(input => from(handler(input)).pipe(map(() => input))) ) return { // 调用next时直接返回Promise,等待本次处理完成 next: async (input: T) => { const result = firstValueFrom(process$) subject.next(input) return result }, // 可根据需要暴露原始subject用于额外订阅 subject } } // 使用示例 const trigger = createAwaitableTrigger(async (input: string) => { // 替换为你的complexFunction逻辑 await new Promise(resolve => setTimeout(resolve, 1000)) console.log('异步处理完成', input) }) // 调用时直接await即可 await trigger.next('测试输入') console.log('订阅侧逻辑已执行完成,继续后续操作')
方案2:自定义AwaitableSubject,支持直接await next方法
如果你需要保留原生Subject的subscribe调用方式,允许直接传入async函数作为订阅回调,可以重写Subject的next方法,收集所有订阅返回的Promise并等待全部执行完成:
import { Subject } from 'rxjs' class AwaitableSubject<T> extends Subject<T> { override next(value: T): Promise<void[]> { const pendingPromises: Promise<void>[] = [] // 遍历所有订阅观察者,收集async回调返回的Promise this.observers.forEach(observer => { const callbackReturn = observer.next?.(value) if (callbackReturn instanceof Promise) { pendingPromises.push(callbackReturn) } }) // 返回所有异步回调的执行结果 return Promise.all(pendingPromises) } } // 使用示例 const mySubject$ = new AwaitableSubject<string>() mySubject$.subscribe(async input => { await complexFunction(input) }) // 直接await next即可等待订阅侧逻辑执行完成 await mySubject$.next('测试输入') console.log('所有订阅逻辑已执行完成')
测试用例修正
使用方案2的AwaitableSubject后,你的测试用例可以修改为以下形式即可正常通过:
import { of, delay, tap, firstValueFrom } from 'rxjs' fdescribe('test', () => { it('should work', async () => { const subject = new AwaitableSubject<string>(); let testValue; subject.subscribe(async inp => { await firstValueFrom(of(inp).pipe(delay(1000), tap(t => testValue = t))) }); await subject.next('text'); expect(testValue).toEqual('text'); }); });
内容的提问来源于stack exchange,提问作者Ádám Bozzay
相关产品推荐
相关产品推荐

