【问题标题】:Switch that doesn't waste work不浪费工作的开关
【发布时间】:2017-05-12 16:40:06
【问题描述】:

我有一个 observable 的 observables,其中内部的 observables 每个都产生一个计算成本高的值。我想要像Switch 这样的行为,但是没有浪费任何工作(SwitchFrugal?):

  • 如果有一个当前的内部 observable 被订阅,我想继续从它接收值直到它完成
    • 不应像Switch那样取消订阅
  • 一旦当前的内部 observable 完成,最新的内部 observable(如果在当前之后有的话)应该被订阅

我一直在努力实现这种行为。使用现有的运营商是否可行?或者这需要通过Observable.Create“从头开始”完成吗?

【问题讨论】:

  • 您要找的不是concatMap()exhaustMap() 吗?
  • @martin exhaust/exhaustMapSwitchFirst 相同,这几乎是我想要的,除了它不满足第二个要点:“一旦当前内部可观察对象完成,应该订阅最新的内部 observable(如果当前有的话)"
  • 你能添加一个你正在寻找的大理石图吗?
  • @DanielT。我添加了一个大理石图。
  • @TimothyShields - 这不正是Merge 所做的吗?

标签: rxjs system.reactive reactive-programming


【解决方案1】:

我会这样做:

const subject = new Subject();
const start = Scheduler.async.now();

// Simulated source Observable
const source = Observable
    .timer(0, 900)
    .map(i => Observable.defer(() => {
        console.log('Subscribe inner Observable:', i);
        return Observable
            .timer(0, 400)
            .take(6)
            .map(j => i.toString() + ':' + j.toString());
    }))
    .share();

source
    .merge(subject
        .withLatestFrom(source)
        .map(([v, observable]) => {
            return observable;
        })
    )
    .exhaustMap(obs => obs.finally(() => subject.next()))
    .subscribe(val => console.log(Scheduler.async.now() - start, val));

我正在模拟一个发出 Observables 的源 Observable。外部 Observable 的发射速度比内部 Observable 完成的速度快,因此这应该模拟您的情况。

然后我将source 合并到链中,但仅在subject 发出时。然后在finally 运算符中触发主题。

输出如下:

Subscribe inner Observable: 0
45 '0:0'
452 '0:1'
856 '0:2'
1260 '0:3'
1665 '0:4'
2065 '0:5'
Subscribe inner Observable: 2
2069 '2:0'
2472 '2:1'
2876 '2:2'
3280 '2:3'
3683 '2:4'
4084 '2:5'
Subscribe inner Observable: 4
4086 '4:0'
4487 '4:1'

请注意,当第一个内部 Observable 完成来自源的最新发射时,是带有 i === 2 的 Observable。如果你运行这段代码,你会发现这三个发射之间没有时间间隔(jsbin现在坏了,所以我无法登录并制作演示):

2065 '0:5'
Subscribe inner Observable: 2
2069 '2:0'

如果您将此与没有 merge() 的默认行为进行比较,您会发现 exhaustMap 需要等到源发出另一个发射:

source
    .exhaustMap(obs => obs.finally(() => subject.next()))
    .subscribe(console.log);

这将打印以下内容。请注意时间间隔,它使用 i === 3 而不是 2 订阅了内部 Observable:

Subscribe inner Observable: 0
45 '0:0'
449 '0:1'
853 '0:2'
1257 '0:3'
1659 '0:4'
2064 '0:5'
Subscribe inner Observable: 3
2748 '3:0'
3151 '3:1'
3553 '3:2'
3953 '3:3'
4355 '3:4'
4759 '3:5'
Subscribe inner Observable: 6
5458 '6:0'
5863 '6:1'

编辑:

为了避免两次订阅同一个内部 Observable(假设这些是冷 Observable),我可以跟踪我已经订阅了哪些 Observable 索引以及接下来需要哪些索引:

我将使源以随机间隔和更少的值发出:

const source = Observable.range(0, 100, Scheduler.async)
    .concatMap(i => Observable.of(i).delay(Math.random() * 3000))
    .map(i => Observable.defer(() => {
        console.log('Subscribe inner Observable:', i);
        return Observable
            .timer(0, 400)
            .take(4)
            .map(j => i.toString() + ':' + j.toString());
    }))
    .map((observable, index) => [observable, index])
    .share();

然后发送我们使用subject.next() 处理的索引并忽略我们不想要的 Observables:

source
    .merge(subject
        .withLatestFrom(source)
        .map(([processedIndex, observableAndIndex]) => {
            let observableIndex = observableAndIndex[1];
            if (processedIndex < observableIndex) {
                return observableAndIndex;
            }
            return false;
        })
        .filter(Boolean)
    )
    .exhaustMap(([observable, index]) => observable.finally(() => subject.next(index)))
    .subscribe(val => console.log(Scheduler.async.now() - start, val));

输出非常相似,但即使之前的 Observable 很快完成,我们也不会再次订阅它(例如 Observables 12 之间的时间间隔):

Subscribe inner Observable: 0
2803 '0:0'
3208 '0:1'
3615 '0:2'
4016 '0:3'
Subscribe inner Observable: 1
4853 '1:0'
5254 '1:1'
5658 '1:2'
6061 '1:3'
Subscribe inner Observable: 2
7814 '2:0'
8218 '2:1'
8622 '2:2'
9026 '2:3'
Subscribe inner Observable: 3
9180 '3:0'
9583 '3:1'
9987 '3:2'
10391 '3:3'
Subscribe inner Observable: 5
10393 '5:0'
10796 '5:1'

【讨论】:

  • 这看起来非常接近,谢谢!我认为唯一的问题是,如果一个内部 observable 完成并且没有更新的内部 observable 到达,那么相同的 observable 将再次被订阅。如果源 observable 发出一个单独的 observable 然后永远不会完成,这不会一遍又一遍地订阅同一个内部 observable 吗? (我要到星期一才能试用。)
  • @TimothyShields 对于冷的 Observables 也是如此(如果您在完成后订阅了一个热的 Observable,它不会做任何事情)。查看我的更新,您可以简单地跟踪您已经订阅的索引,并最终忽略您不想要的 Observables。
猜你喜欢
  • 2023-04-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-11-05
  • 2010-10-23
  • 2014-06-11
  • 2014-10-20
相关资源
最近更新 更多