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

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);
  }
}

核心逻辑解释

  1. 多次Resolve支持:
    构造函数中传递给setup的resolve函数本质是调用BehaviorSubject.next(),每次调用都会更新Subject保存的最新值,天然支持多次触发。

  2. Await返回最后值:
    BehaviorSubject会始终保存最新发射的值,当调用await时,then方法内的订阅会立刻获取到当前最新的非初始值(如果已有),并通过take(1)完成Promise,直接返回该值。

  3. 未Resolve时挂起:
    如果Subject处于初始的INITIAL_STATE,filter操作符会过滤掉该状态,订阅会一直等待直到第一个真正的resolve值被发射,Promise才会完成,和标准Promise的挂起行为一致。

  4. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 01:45:06