RxJS新手求助:如何实现可多次resolve、带reset的类Promise API?
用RxJS实现类Promise的多Resolve API
你可以基于RxJS的BehaviorSubject来实现这个需求,它天然支持保存最新值、订阅后立刻获取最新值(或等待第一个值),再包装成支持await的类即可满足所有要求。
完整实现代码
import { BehaviorSubject, filter, take, Subscription } from 'rxjs'; // 用唯一Symbol标记初始未resolve状态,避免和业务值冲突 const INITIAL_STATE = Symbol('initial-unresolved'); class MyObservableApi<T> { private subject: BehaviorSubject<T | typeof INITIAL_STATE>; private sourceSubscription?: Subscription; constructor(setup: (resolve: (value: T) => void) => Subscription) { // 初始化BehaviorSubject为未resolve状态 this.subject = new BehaviorSubject<T | typeof INITIAL_STATE>(INITIAL_STATE); // 传递给setup的resolve函数,本质是发射新值到Subject const resolve = (value: T) => { this.subject.next(value); }; // 保存setup返回的订阅,用于reset时取消 this.sourceSubscription = setup(resolve); } // 实现then方法,让类支持await语法 then<TResult1 = T, TResult2 = never>( onFulfilled?: ((value: T) => TResult1 | PromiseLike<TResult1>) | null | undefined, onRejected?: ((reason: any) => TResult2 | PromiseLike<TResult2>) | null | undefined ): Promise<TResult1 | TResult2> { return new Promise((resolve, reject) => { const subscription = this.subject.pipe( // 过滤掉初始未resolve状态,只处理真正的resolve值 filter(value => value !== INITIAL_STATE), // 只取第一个值(即当前最新值),完成Promise take(1) ).subscribe({ next: (value) => { try { const result = onFulfilled ? onFulfilled(value as T) : value; resolve(result as TResult1); } catch (err) { reject(err); } }, error: (err) => { try { const result = onRejected ? onRejected(err) : Promise.reject(err); resolve(result as TResult2); } catch (innerErr) { reject(innerErr); } } }); }); } // 实现catch方法,对齐Promise语法 catch<TResult = never>( onRejected?: ((reason: any) => TResult | PromiseLike<TResult>) | null | undefined ): Promise<T | TResult> { return this.then(undefined, onRejected); } reset() { // 取消原始数据源的订阅,防止内存泄漏 this.sourceSubscription?.unsubscribe(); // 将Subject重置为未resolve状态,后续await会重新等待新值 this.subject.next(INITIAL_STATE); } }
核心逻辑解释
多次Resolve支持:
构造函数中传递给setup的resolve函数本质是调用BehaviorSubject.next(),每次调用都会更新Subject保存的最新值,天然支持多次触发。Await返回最后值:
BehaviorSubject会始终保存最新发射的值,当调用await时,then方法内的订阅会立刻获取到当前最新的非初始值(如果已有),并通过take(1)完成Promise,直接返回该值。未Resolve时挂起:
如果Subject处于初始的INITIAL_STATE,filter操作符会过滤掉该状态,订阅会一直等待直到第一个真正的resolve值被发射,Promise才会完成,和标准Promise的挂起行为一致。Reset方法:
- 取消原始数据源的订阅(比如你伪代码中的
subscribe返回的订阅),避免内存泄漏; - 将Subject重置为初始状态,后续的
await或then会重新等待新的resolve调用。
- 取消原始数据源的订阅(比如你伪代码中的
你的伪代码适配示例
interface Auth { token: string; } // 模拟你的订阅逻辑,返回RxJS Subscription用于取消 function subscribe(selector: (state: any) => Auth, callback: (auth: Auth) => void): Subscription { // 模拟定时发射auth更新 const intervalId = setInterval(() => { callback({ token: `auth-token-${Date.now()}` }); }, 1000); return new Subscription(() => clearInterval(intervalId)); } // 初始化你的自定义API const myAuthObservable = new MyObservableApi<Auth>((resolve) => subscribe( (s) => s.auth, (auth) => { if (auth) { resolve(auth); // 可以多次调用 } } ) ); // 模拟事件触发 const emitter = { on(event: string, handler: () => void) { setTimeout(handler, 3000); // 3秒后触发start事件 } }; emitter.on("start", async () => { // 拿到最后一次resolve的auth值,若未resolve则等待 const auth = await myAuthObservable; console.log("当前Auth:", auth); // 重置API,后续await会等待新的resolve myAuthObservable.reset(); // 重置后再次await,会等待新的auth值 const newAuth = await myAuthObservable; console.log("新的Auth:", newAuth); });
内容的提问来源于stack exchange,提问作者nikitakot
相关产品推荐
相关产品推荐

