【发布时间】:2019-12-29 01:59:22
【问题描述】:
我想使用现有的 rxjs 操作符来创建一个多播 observable,每次订阅它时都会重新订阅它的源,但在源完成时不会取消订阅它的订阅者。
目的是混合使用共享同一数据源的短期组件和长期组件。当创建新的短期组件时,应该为新组件和任何同时订阅数据源的长期组件更新共享数据。
使用 Observable.create(),我能够创建一个自定义的 observable 来产生这种行为,但感觉可能有一个开箱即用的解决方案,不需要编写自定义代码。
这是我尝试过的。
import { of, Observable, Observer, zip, interval, merge } from "rxjs";
import { filter, map, publish, tap, delay, shareReplay, share } from "rxjs/operators";
class MySubject {
constructor(private observable: Observable<any>) {}
sourceActive = false;
subscribers: Array<Observer<any>> = [];
public subscribe(observer: Observer<any>) {
this.subscribers.push(observer);
if (!this.sourceActive) {
console.log("subscribing");
this.sourceActive = true;
this.observable.subscribe(
x => this.subscribers.forEach(sub => sub.closed || sub.next(x)),
x => this.subscribers.forEach(sub => sub.closed || sub.error(x)),
() => this.sourceActive = false
);
}
}
}
const source$ = of(1).pipe(
tap(x=>console.log("invoked cold")),
delay(2000)
);
const mySubject = new MySubject(source$);
const super$ = Observable.create(observer => mySubject.subscribe(observer));
const sub1 = super$.subscribe(x => console.log("sub 1"));
const sub2 = super$.subscribe(x => console.log("sub 2"));
setTimeout(x => {
sub1.unsubscribe();
super$.subscribe(x => console.log("sub 3"));
}, 3000);
【问题讨论】:
标签: angular typescript rxjs rxjs-pipeable-operators