【问题标题】:Retrieve Observable subscribers and make them subscribe to another Observable检索 Observable 订阅者并让他们订阅另一个 Observable
【发布时间】:2019-12-03 21:41:05
【问题描述】:

简单地说

给定一个现有的 Observable(尚未完成),有没有办法检索关联的订阅者(传递给 subscribe 的函数)以使他们订阅另一个 Observable?

上下文

我的应用程序中的服务帮助创建 SeverEvent 连接,返回一个 ConnectableObservable 到代理连接,并允许使用 publish 运算符进行多播。该服务通过内部存储跟踪现有连接:

store: {[key: string]: ConnectionTracker};

// …

interface ConnectionTracker {
    url: string;
    eventSource: EventSource;
    observable: rx.ConnectableObservable<any>;
    subscription: rx.Subscription;
    observer: rx.Observer<any>;
    data?: any; // Arbitrary data
}

在创建连接后,如果关联的跟踪器已经存在(使用连接的端点生成身份),服务应该:

  • ok 关闭现有tracker的ServerEvent连接
  • ok 打开一个新的 SerevrEvent 连接(因此是一个新的 ConnectableObservable)
  • 用新的 observable 替换现有跟踪器的 Observable 但现在让现有订阅者改为订阅新的 Observable

这是创建ConnectionTrackers

的代码部分
/**
* Create/Update a ServerEvent connection tracker
*/
createTracker<T>(endpoint: string, queryString: string = null): ConnectionTracker
{
    let fullUri = endpoint + (queryString ? `?${queryString}` : '')
        , tracker = this.findTrackerByEndpoint(endpoint) || {
            observable: null,
            fullUri: fullUri,
            eventSource: null,
            observer: null,
            subscription: null
        }
    ;

    // Tracker exists
    if (tracker.observable !== null) {
        // If fullUri hasn't changed, use the tracker as is
        if (tracker.fullUri === fullUri) {
            return tracker;
        }

        // At this point, we know "fullUri" has changed, the tracker's
        // connection should be replaced with a fresh one

// ⇒ TODO
// ⇒ Gather old tracker.observable's subscribers/subscriptions to make
//   them subscribe to the new Observable instead (created down below)

        // Terminate previous connection and clean related resouces
        tracker.observer.complete();
        tracker.eventSource.close();
    }

    tracker.eventSource = new EventSource(<any>fullUri, {withCredentials: true});
    tracker.observable = rx.Observable.create((observer: rx.Observer<T>) => {
            // Executed once
            tracker.eventSource.onmessage = e => observer.next(JSON.parse(e.data));
            tracker.eventSource.onerror = e => observer.error(e);
            // Keep track of the observer
            tracker.observer = observer;
        })
        // Transform Observable into a ConnectableObservable for multicast
        .publish()
    ;

    // Start emitting right away and also keep a reference to 
    // proxy subscription for later disposal
    tracker.subscription = tracker.observable.connect();

    return tracker;
}

谢谢。

【问题讨论】:

  • 看起来你可以只使用 switchMap 来返回“新的” Observable。订阅者将保持原样,但会从这个“新” Observable 接收值。
  • 我担心 switchMap 不是我想要的。我的意图是创建一个替代 Observable,但从以前的 Observable 中恢复注册的观察者。一旦 ServerEvent 连接关闭,相关的 Observable 就会过时(不再有源 → 没有 next() 调用),可以删除对 observable 的引用。 switchMap 只构建一个 Observable 链:理论上原始的 Observable 仍然在值班(据我所知),但是由于相关的连接已关闭,因此不会再发出任何值,并且新的 Observable 不会有有机会接管。

标签: rxjs rxjs-observables


【解决方案1】:

与其尝试手动将订阅者从一个 Observable 转移到另一个 Observable,不如为侦听器提供一个 Observable,它会在需要时自动切换到另一个 Observable。

您可以通过使用 高阶 Observable(一个发出 Observables 的 Observable)来做到这一点,该 Observable 总是切换到最新的内部 Observable。

基本概念

// a BehaviorSubject is used so that late subscribers also immediately get the most recent inner Observable
const higherOrderObservable = new BehaviorSubject<Observable<any>>(EMPTY);

// pass new Observable to listeners
higherOrderObservable.next(new Observable(..));

// get most recent inner Observable
const currentObservable = higherOrderObservable.pipe(switchMap(obs => obs));
currentObservable.subscribe(valueFromInnerObservable => { .. })

你的情况

为每个端点创建一个BehaviorSubject (tracker supplier),它发出当前应该用于该目的的 Observable (tracker) 端点。当应该为给定的 endpoint 使用不同的 tracker 时,将这个新的 Observable 传递给 BehaviorSubject。让您的听众订阅BehaviorSubject (tracker supplier),它会自动为他们提供正确的tracker,即切换到当前应该使用的 Observable。

您的代码的简化版本可能如下所示。具体情况取决于您如何在整个应用中使用函数 createTracker

interface ConnectionTracker {
  fullUri: string;
  tracker$: ConnectableObservable<any>;
}

// Map an endpoint to a tracker supplier.
// This is your higher order Observable as it emits objects that wrap an Observable
store: { [key: string]: BehaviorSubject<ConnectionTracker> };
closeAllTrackers$ = new Subject();

// Creates a new tracker if necessary and returns a ConnectedObservable for that tracker. 
// The ConnectedObservable will always resemble the current tracker.
createTracker<T>(endpoint: string, queryString: string = null): Observable<any> {
  const fullUri = endpoint + (queryString ? `?${queryString}` : '');
  // if no tracker supplier for the endpoint exists, create one
  if (!store[endpoint]) {
    store[endpoint] = new BehaviorSubject<ConnectionTracker>(null);
  }
  const currentTracker = store[endpoint].getValue();

  // if no tracker exists or the current one is obsolete, create a new one
  if (!currentTracker || currentTracker.fullUri !== fullUri) {
    const tracker$ = new Observable<T>(subscriber => {
      const source = new EventSource(fullUri, { withCredentials: true });
      source.onmessage = e => subscriber.next(JSON.parse(e.data));
      source.onerror = e => subscriber.error(e);
      return () => source.close(); // on unsubscribe close the source
    }).pipe(publish()) as ConnectableObservable<any>;
    tracker$.connect();
    // pass the new tracker to the tracker supplier
    store[endpoint].next({ fullUri, tracker$ });
  }
  // return the tracker supplier for the given endpoint that always switches to the current tracker
  return store[endpoint].pipe(
    switchMap(tracker => tracker ? tracker.tracker$ : EMPTY), // switchMap will unsubscribe from the previous tracker and thus close the connection if a new tracker comes in
    takeUntil(this.closeAllTrackers$) // complete the tracker supplier on emit
  );
}

// close all trackers and remove the tracker suppliers
closeAllTrackers() {
  this.closeAllTrackers$.next();
  this.store = {};
}

如果您想立即关闭所有跟踪器连接并且现有订阅者应该收到complete 通知,请致电closeAllTrackers。 如果您只想关闭一些跟踪器连接但不希望现有订阅者收到complete 通知,以便他们继续侦听将来提供的新跟踪器,请为每个跟踪器调用store[trackerEndpoint].next(null)

【讨论】:

  • 如果我需要一次性删除所有现有连接怎么办?循环遍历 store 自己的属性并在每个 BehaviorSubject 上调用 .complete() 就足够了吗?
  • 您可能不想自己完成 BehaviorSubjects,因为如果您愿意,之后您将无法发出任何新的跟踪器。相反,只需将null 传递给每个跟踪器供应商:store[endpoint].next(null)。我编辑了最后一个 return 语句,这样如果 null 被传递给一个 BehaviorSubject 它将切换到一个刚刚完成的 EMPTY observable。这样,您之前的跟踪器将被取消订阅,并且在将新的 Observable 传递给 BehaviorSubject 之前,侦听器将不会收到任何数据。
【解决方案2】:

如果你试图做一些事情,比如将订阅者移动到不同的 observable,那么你只是没有做 RxJS 中想要做的事情。任何此类操纵基本上都是黑客行为。

如果您偶尔产生一个新的 observable(例如通过发出请求),并且您希望某些订阅者始终订阅其中最新的,那么解决方案如下:

  private observables: Subject<Observable<Data>> = new Subject();

  getData(): Observable<Data> {
    return this.observables.pipe(switchAll());
  }

  onMakingNewRequest(newObservable: Observable<Data>) {
    this.observables.push(newObservable);
  }

通过这种方式,您可以公开客户端订阅的单个 observable(通过 getData()),但通过推送到 this.observables,您可以更改用户看到的实际数据源。

至于关闭连接和类似的东西,你的 observable(每个请求或其他东西创建的那个)基本上应该在取消订阅时负责释放和关闭东西,那么你不需要做任何额外的处理,从您推送新的那一刻起,以前的 observable 将自动取消订阅。详细信息取决于您联系的实际后端。

【讨论】:

  • 你的第二句话听起来像我正在寻找的东西,我觉得 Subject / switchAll 组合正是我所需要的。我会挖掘这个,谢谢。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-08-29
  • 1970-01-01
  • 1970-01-01
  • 2019-12-29
  • 1970-01-01
相关资源
最近更新 更多