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

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 15:06:03