【问题标题】:Subscribe to observable A when someone subscribe to observable B当有人订阅 observable B 时订阅 observable A
【发布时间】:2021-11-26 08:26:04
【问题描述】:

我有一个 API 调用是可观察的情况。 我还有一个可观察的 B 来跟踪该查询的进度:

  generateUrl$(uploadConfig: GenerateUploadUrlVariables, file: File) {
    const generation$ = new BehaviorSubject<ActionWithResult<{ url: string }>>({
      state: 'INIT',
      result: null,
    });

    const complete = function (result: ActionWithResult<{ url: string }>) {
      generation$.next(result);
      generation$.complete();
    };

    this.generateUploadUrl(uploadConfig).pipe(
      switchMap((result) => {
        generation$.next({ state: 'UPLOADING', result: null });
        const url = result.data.generateUploadUrl.url || '';
        return this.httpClient
          .pipe(
            tap(() => {
              complete({ state: 'SUCCEEDED', result: { url } });
            }),
            catchError((e) => {
              complete({ state: 'FAILED', result: null });
              return throwError(() => e);
            })
          );
      }),
      catchError((e) => {
        complete({ state: ActionState.FAILED, result: null });
        return throwError(() => e);
      })
    ).subscribe();

    return generation$
  }

现在可以了

this.generateUrl$(option, file).subscribe(e=&gt; {console.log(e)})

我会得到查询的状态,以及查询成功后的结果。

但问题是,如果有人犯了错误,干脆做:

this.generateUrl$()没有订阅,不会跟踪结果,但API仍然会被调用。

我想将 api observable this.generateUploadUrl 与跟踪 observable generation$ 绑定

类似

return generation$.pipe(
    whenSomeoneSubscribeToIt(() => {
       this.generateUploadUrl(uploadConfig)
       .pipe(//all the stuff above)
       .subscribe()
    })
)

有可能吗?

【问题讨论】:

    标签: rxjs


    【解决方案1】:

    天真的解决方案

    我创建了一个自定义运算符,它可以在订阅流之后但在接收到它的第一次发射之前运行效果。我称它为initialize,但这几乎就是你用whenSomeoneSubscribeToIt 描述的内容。

    这里是:

    export function initialize<T>(
      effect: () => void
    ): MonoTypeOperatorFunction<T> {
    
      return s => new Observable(ob => {
        const sub = s.subscribe(ob);
        effect();
        return sub;
      });
    
    }
    

    现在您可以让您的订阅生效了:

    generateUrl$(
      uploadConfig: GenerateUploadUrlVariables, 
      file: File
    ) {
    
      // Some stuff here
    
      return generation$.pipe(
        initialize(() => {
          this.generateUploadUrl(uploadConfig).pipe(
            // some more stuff here
          ).subscribe();
        })
      );
    
    }
    

    惯用的 RxJS

    您根本不需要上面的BehaviourSubject。您基本上是在使用它来强制生成流。

    正如另一个答案所述,您可以直接向观察者创建相同的流。这种方法要好一些,但仍然会遇到使用行为主题的许多问题。例如,该答案会创建一个无法保证完成的可观察对象,并且在您取消订阅时不会管理它的内部订阅。

    我会说惯用的 RxJS 方法是在使用 RxJS 运算符之后以声明方式创建您要的流。

    generateUrl$(
      uploadConfig: GenerateUploadUrlVariables, 
      file: File
    ): Observable<StateWithResult> {
    
      return this.generateUploadUrl(uploadConfig).pipe(
        map(result => result.data.generateUploadUrl.url || ''),
        switchMap(url => this.httpClient.pipe(
          ignoreElements(),
          startWith({ state: 'UPLOADING', result: null }),
          concatWith(of({ state: 'SUCCEEDED', result: { url } }))
        )),
        catchError(_ => of({ state: 'FAILED', result: null })),
        startWith({ state: 'INIT', result: null })
      );
    
    }
    

    我在这里所做的一个改变是我不会重新抛出你的错误。看起来你无论如何都没有在下游使用它们,所以这是一种更清洁的处理它们的方法。你可以(当然!)回去重新扔掉它们。

    【讨论】:

    • 有趣的答案!类似于我想出的。我从来不知道ignoreElements,所以今天早上这对我来说是一个很好的小发现:-)
    • @BizzyBob。 :) ignoreElements 很好,因为它告诉 typescript 生成的 observable 没有 next 排放。否则与filter(_ =&gt; false) 相同。
    【解决方案2】:

    这可能是在这种情况下不使用像BehaviorSubject 这样的多播的原因之一。使用 new Observable() 函数创建一个 observable 会更适合您。

    请注意,我还调整了其他一些实现细节,例如使用 tap 运算符的 finalizeerror 回调而不是内部 catchError 运算符。

    试试下面的

    import { Observable, Observer, throwError } from 'rxjs';
    import { catchError, finalize, switchMap, tap } from 'rxjs/operators';
    
    generateUrl$(uploadConfig: GenerateUploadUrlVariables, file: File): Observable<any> {
      return new Observable((generation$: Observer) => {
        generation$.next({ state: 'INIT', result: null });
        this.generateUploadUrl(uploadConfig).pipe(
          switchMap((result: any) => {
            generation$.next({ state: 'UPLOADING', result: null });
            const url = result.data.generateUploadUrl.url || '';
            return this.httpClient.get(someUrl).pipe(
              tap({
                next: (res: any) => generation$.next({ state: 'SUCCEEDED', result: { url } }),
                error: (error: any) => generation$.next({ state: 'FAILED', result: null })
              }),
              finalize(() => generation$.complete())
            );
          }),
          catchError((error: any) => {
            generation$.next({ state: ActionState.FAILED, result: null });
            generation$.complete();
            return throwError(() => error);
          })
        ).subscribe();
      });
    }
    

    【讨论】:

      【解决方案3】:

      即使消费者没有调用.subscribe(),http 调用也会触发的原因是因为您在函数内部的generateUploadUrl 上调用.subscribe()

      您可以通过简单地返回 generateUploadUrl 返回的 observable 并将其结果传递给您想要的值来获得您想要的行为(我想说的是正确的反应行为)。我们可以使用startWith 来发出您的初始值:

        generateUrl$(uploadConfig: GenerateUploadUrlVariables, file: File): Observable<ActionWithResult> {
      
          return this.generateUploadUrl(uploadConfig).pipe(
            map(result => result.data.generateUploadUrl.url || ''),
            switchMap(url => this.httpClient.pipe(
              map(() => ({ state: ActionState.Succeeded, result: { url } })),
              startWith({ state: ActionState.Uploading, result: null })
            )),
            catchError(() => of({ state: ActionState.Failed, result: null })),
            startWith({ state: ActionState.Init, result: null })
          );
        }
      

      这是一个有效的 StackBlitz 演示。

      请注意,您没有在函数内部订阅,这首先导致了您的不良行为。

      如果您需要多播行为(如果您有多个订阅者),您只需在管道末尾添加 shareReplay

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-12-29
        • 2020-08-29
        • 2023-04-02
        • 2018-08-28
        • 2017-03-02
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多