在RxJS 6中重置ReplaySubject

NRa*_*Raf 5 javascript rxjs

我有一个可过滤的“活动日志”,当前使用来实现ReplaySubject(因为一些组件在使用它,并且它们可能在不同的时间进行订阅)。

当用户更改过滤器设置时,会发出一个新请求,但是结果将附加到ReplaySubject而不是替换它。

我想知道是否有ReplaySubject某种方法可以更新,使其仅使用诸如的内容通过新项目发送switchMap

否则,我可能需要使用BehaviorSubject返回所有活动条目的数组的,或者重新创建ReplaySubject并通知用户(可能通过使用另一个可观察的对象)来取消订阅并重新订阅新的可观察对象。

car*_*ant 10

如果您希望能够在不让其订阅者明确取消订阅和重新订阅的情况下重置主题,则可以执行以下操作:

import { Observable, Subject } from "rxjs";
import { startWith, switchMap } from "rxjs/operators";

function resettable<T>(factory: () => Subject<T>): {
  observable: Observable<T>,
  reset(): void,
  subject: Subject<T>
} {
  const resetter = new Subject<any>();
  const source = new Subject<T>();
  let destination = factory();
  let subscription = source.subscribe(destination);
  return {
    observable: resetter.asObservable().pipe(
      startWith(null),
      switchMap(() => destination)
    ),
    reset: () => {
      subscription.unsubscribe();
      destination = factory();
      subscription = source.subscribe(destination);
      resetter.next();
    },
    subject: source
  };
}
Run Code Online (Sandbox Code Playgroud)

resettable 将返回一个包含以下内容的对象:

  • 一个observable到哪些订户到重新设定的主题应订阅;
  • 一个subject在其中你会打电话来nexterrorcomplete; 和
  • 一个reset将重置(内)对象的功能。

您可以这样使用它:

import { ReplaySubject } from "rxjs";
const { observable, reset, subject } = resettable(() => new ReplaySubject(3));
observable.subscribe(value => console.log(`a${value}`)); // a1, a2, a3, a4, a5, a6
subject.next(1);
subject.next(2);
subject.next(3);
subject.next(4);
observable.subscribe(value => console.log(`b${value}`)); // b2, b3, b4, b5, b6
reset();
observable.subscribe(value => console.log(`c${value}`)); // c5, c6
subject.next(5);
subject.next(6);
Run Code Online (Sandbox Code Playgroud)


Ere*_*z.S 5

这是一个使用之前发布的可重置工厂的类,因此您可以使用 const myReplaySubject = new ResettableReplaySubject<myType>()

import { ReplaySubject, Subject, Observable, SchedulerLike } from "rxjs";
import { startWith, switchMap } from "rxjs/operators";

export class ResettableReplaySubject<T> extends ReplaySubject<T> {

reset: () => void;

constructor(bufferSize?: number, windowTime?: number, scheduler?: SchedulerLike) {
    super(bufferSize, windowTime, scheduler);
    const resetable = this.resettable(() => new ReplaySubject<T>(bufferSize, windowTime, scheduler));

    Object.keys(resetable.subject).forEach(key => {
        this[key] = resetable.subject[key];
    })

    Object.keys(resetable.observable).forEach(key => {
        this[key] = resetable.observable[key];
    })

    this.reset = resetable.reset;
}


private resettable<T>(factory: () => Subject<T>): {
    observable: Observable<T>,
    reset(): void,
    subject: Subject<T>,
} {
    const resetter = new Subject<any>();
    const source = new Subject<T>();
    let destination = factory();
    let subscription = source.subscribe(destination);
    return {
        observable: resetter.asObservable().pipe(
            startWith(null),
            switchMap(() => destination)
        ) as Observable<T>,
        reset: () => {
            subscription.unsubscribe();
            destination = factory();
            subscription = source.subscribe(destination);
            resetter.next();
        },
        subject: source,
    };
}
}
Run Code Online (Sandbox Code Playgroud)