【问题标题】:How to send to a Subject different Observable results without having to break the observable flow subscribing?如何在不中断可观察流订阅的情况下向主题发送不同的可观察结果?
【发布时间】:2020-08-11 08:41:43
【问题描述】:

用一个简洁的例子编辑了这个问题。我想这会让人们更容易理解。注意这个例子是超级简化的,通常我在不同的组件之间分层,但是对于这个问题就足够了。

拿这个组件。它采用获取的对象的名称和获取下一个对象的按钮。为了获得下一个 REST 请求的值,我知道除了订阅答案之外别无他法,我想要的是像“combineLatest”这样的东西,但对于“未来”,这样我就可以组合最新的流.

import { Component, VERSION, OnInit } from '@angular/core';
import { BehaviorSubject, Observable } from 'rxjs';
import { HttpClient } from '@angular/common/http';

@Component({
  selector: 'my-app',
  templateUrl: './app.component.html',
  styleUrls: ['./app.component.css']
})
export class AppComponent implements OnInit {
  
    private readonly PEOPLE_API_ENDPOINT = `https://swapi.dev/api/people/`;

   private characterSubject : BehaviorSubject<any> = new BehaviorSubject<any>({name: 'loading'});
  private currentCharacter: number = 1;
  
  character$ : Observable<any> = this.characterSubject.asObservable();
  
  constructor(
    private http: HttpClient
  ) {}
    
  ngOnInit() {
    this.updateCurrentCharacter();
  }

  nextCharacter() : void {
    this.currentCharacter ++;
    this.updateCurrentCharacter();
  }

  //I would want to avoid subscribing, instead
  //I would like some sort of operation to send the stream
  //emissions to the subject. As to not break the observable
  //chain up until presentation, like the best practices say.
  private updateCurrentCharacter() : void {
    this.fetchCharacter(this.currentCharacter)
      .subscribe( 
        character => this.characterSubject.next(character)
      );
  }

  private fetchCharacter (id: number) : Observable<any> {
        return this.http.get(this.PEOPLE_API_ENDPOINT + `${id}/`);
  }
}
<span>{{ character$ | async }} </span>

<button (click)="nextCharacter()">Next character</button>

Online demo

有没有办法做到这一点?做类似“emitIn(characterSubject)”的事情。我认为没有像这样动态地将源排放添加到源。

【问题讨论】:

  • hm 我不认为我明白 - 视图应该从 currentDatadataRepository 获取数据吗?
  • 如果您想延迟实际订阅时间,您可以使用publishconnect。这是您要找的东西吗?
  • 如果你经常遇到这个问题并且不想完全使用像 NgRx 这样的状态管理库,你可能想研究一下纯 RxJS 的状态管理。 Video - Article
  • 我想我会用一个具体的例子来编辑这个问题。也许会更容易理解。

标签: angular rxjs observable reactive subject


【解决方案1】:

如果我理解正确,您的服务具有触发 http 调用的某种方法,例如 dataRepository.fetch(id),并且您有不同的组件需要在调用响应到达时做出反应。

如果是这种情况,有很多方法可以处理这种要求,一种是使用作为服务的公共属性公开的主题,这是我理解你想要做的。

要实现此行为,您编写的代码就是您所需要的,换句话说,这是可以的

dataRepository.fetch(id)
  .subscribe(
    newData => currentData.next(newData)
  )

如果你想让它更完整并管理errorcomplete的情况,你可以这样写

dataRepository.fetch(id)
  .subscribe(currentData)

在最后一种形式中,您将currentData 作为subscribe 的观察者参数传递。 currentDatathough 也是一个 Observable,因此可以 nexterrorcomplete

如果您使用ReplaySubject,您可以添加存储最后结果并将它们呈现给在结果通知后创建的组件的可能性。

【讨论】:

  • 那么除了必须与可观察对象同步之外,别无他法? (使用订阅)。耻辱D:。虽然我有点明白,如果你订阅了一个 observable,你就会得到 observable 的所有优点,所以它会自己处理这 3 个操作。但是订阅不会破坏反应链吗?喜欢 || HTTP 获取 -> 反应式/异步 ||订阅 -> 命令式/同步 ||服从异步管道 -> 反应式/异步 ||中间的中断是我想要避免的。但从你说的情况来看,似乎别无选择。
  • 顺便说一句,我没有得到你的最后一个例子。因为事情是每次事件触发时,我都会创建一个新的 fetch,但是组件不知道有一个新的源,我必须手动更新它们中的每一个,这有点违背了 observables 的目的跨度>
  • @JoshiRaez 您有 1 个事件要多路复用到 n 个通知。所以我认为你对 RxJS 没有任何其他方法。另外,我没有看到 Reactive 是异步的,而 Imperative 是同步的。 RxJS 可以是完美的反应式和完美的同步。例如,考虑一个对来自 WebSockets 的消息做出反应的应用程序。反应可以是完全同步的,它仍然是一种反应式的编码方式。 This article 提供更多详细信息。
  • @JoshiRaez 关于您的最后一条评论,您是对的。 share 运算符不适用于您的情况。我已编辑我的回复并将其删除。
【解决方案2】:

检索到的原始组件 observable 是动态更改的,但这真的很危险(任何管道都可能破坏任何 observable)。我需要另一个不叫管道的东西。

我不同意。您可以使用基于您的输入数据返回可观察对象的函数轻松创建动态可观察对象,并且不会中断。

function getUserFriendsById(userId: string): Observable<User[]>{
return service.getProfileById(userId).pipe(
  mergeMap(user=>{
    return getUsersByArrayOfIds(user.friends);
  })
 )
}

getUserFriendsById('exampleUserId').subscribe(friends=>{...});

但如果您认为如果内部出现问题,您可以简单地使用 CatchError。

function getUserFriendsById(userId: string): Observable<User[]>{
return service.getProfileById(userId).pipe(
  mergeMap(user=>{
    return getUsersByArrayOfIds(user.friends).pipe(
      catchError(err => {
      return of([]) // will return empty friend list if this APi call fails.
    }))
  }),
  catchError(err=>{
   return of([]); // will return empty friend list if this APi call fails.
  })
 )
}

getUserFriendsById('exampleUserId').subscribe(friends=>{...});

我们使用mergeMap只是因为我们需要第一个Observable的数据来发送第二个请求。

如果你正在寻找的东西是从基础中混合这些可观察的东西,那么这将使所有 rxjs 运算符无用,因为这就是它们的目的。

但是如果您不想将数据发送到您的 currentData Observable 并且您知道您的数据源,那么您甚至不需要为它创建单独的 Observable .只需使用管道获取数据就足够了。

currentData = dataRepository.fetch(id);
// at some point when you need currentData data you will subscribe to it.

// is same as 

dataRepository.fetch(id)
  .subscribe(
    newData => currentData.next(newData)
  )

甚至当您需要通过dataRepository.fetch(id) 响应调用另一个 Observable 时

dataRepository.fetch(id).pipe(
mergeMap(fetchedId=>{
    return getUserById(fetchedId);
  })
)

但是,如果您想更深入地创建流程而不传递数据,我认为您会感到失望,因为这就是 Observables 在一般情况下的工作方式。而且我想这种方法非常有限且难以维护,因为您最终会到达需要向多个 API 发送不同请求并收集数据的地步,这时流程将中断。 https://www.learnrxjs.io/learn-rxjs/operators/creation/create

// RxJS v6+
import { Observable } from 'rxjs';
/*
  Create an observable that emits 'Hello' and 'World' on  
  subscription.
*/
const hello = Observable.create(function(observer) {
  observer.next('Hello');
  observer.next('World');
  observer.complete();
});

希望对您有所帮助,但如果不是您要找的,请告诉我。

【讨论】:

  • 我想我用“动态”的东西错误地解释了自己。当我的意思是“动态”时,我的意思是这样的原始行为:Source will emit: 1(t0), 2(t2), 3(t3), C First observable from source pipeline with First.那个 observable 将是 1(t0), C 第二个从源管道可观察到的 Map x => x*2。那个 observable 会带来 2(t0), 4(t2), 6(t3), C 注意第一个和第二个管道源(两个更小的 observables)之间是如何不交互的。
  • 但如果它们是动态的,我的意思是每个 """"pipe"""" 都会改变其他 observables,基本上会改变源。我知道,这根本不是响应式的工作方式,我的意思是这也很危险,因为它会将一个可观察的对象耦合到每个管道实例。 Source 将发出:1(t0), 2(t2), 3(t3), C First 可从具有 First 的源管道观察到。那个 observable 将是 1(t0), C 但源也将是 1(t0), C 第二个 observable 从源管道与 Map x => x*2。那个 observable、previous 和 source 现在会变成 2(t0), C
猜你喜欢
  • 2018-06-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-11-26
  • 1970-01-01
  • 1970-01-01
  • 2020-07-21
  • 1970-01-01
相关资源
最近更新 更多