【问题标题】:Can an RxJS 5 Observable source be stopped by another down the chain?一个 RxJS 5 Observable 源可以被另一个链下的源停止吗?
【发布时间】:2016-07-29 08:18:23
【问题描述】:

RxJS 5Angular 2 RC4Typescript 1.9写的应用:

我有两个可观察的链。我想,如果第二个条件满足,第一个立即完成。我的努力似乎不必要地复杂。在下面的示例中,我尝试在第一个 observable 发出 3 个值后停止它:

source = Observable.interval(1000)
        .do(()=>this.print("*******EMITTING from Source*******"))
        .switchMap(count => {
            if(count<3){ //just pass along the value
                return Observable.create(observer=>{
                    observer.next(count);observer.complete()
                })                
            }
            else{ //abort by issuing a non-productive observable
                return Observable.create(observer=>
                    observer.complete()
                )
            }
        })
this.source.subscribe(count=>this.print('Ouput is '+count);

这是输出:

*******EMITTING from Source*******
Output is 0
*******EMITTING from Source*******
Output is 1
*******EMITTING from Source*******
Output is 2
*******EMITTING from Source*******
*******EMITTING from Source*******
*******EMITTING from Source*******

所以,在功能上我得到了我想要的结果,因为更广泛的脚本在三个输出后停止接收通知。但是,我确信有更好的方法。我的问题是:

  1. 上游的 observable 会一直持续发射。我怎样才能阻止它?
  2. 我正在创建一个新的 observable 在每次发射的链条上。我不应该只传递前 3 个值,但在第 4 天中止或完成链吗?

【问题讨论】:

    标签: rxjs observable


    【解决方案1】:

    您可以使用take operator 来执行此操作。take 获取前 N 个事件并完成流。

    this.source = Observable.interval(1000)
      .do(()=>this.print("*******EMITTING from Source*******"))
      .take(3);
    this.source.subscribe(count=>this.print('Ouput is '+count);
    

    您的示例流未完成,因为switchMap 的外部流在内部流完成时未完成。 switchMap() 等于 map().switch()。在您的示例中,map 部分发出 observables,例如:

    1. 下一个(0),完成()
    2. 下一个(1),完成()
    3. 下一个(2),完成()
    4. 完成()
    5. 完成()
    6. 完成()
    7. 完成()
    8. ...(无限继续)...

    switch 部分会切换这些可观察对象并继续等待即将到来的可观察对象。

    编辑

    你的例子也可以写成:

    source = Observable.interval(1000)
            .do(()=>this.print("*******EMITTING from Source*******"))
            .takeWhile(count => count < 3);
    

    编辑 2

    关于您的评论,如果您想在内部流发出 true 时终止流:

    source = Observable.interval(1000)
            .do(()=>this.print("*******EMITTING from Source*******"))
            .switchMap(count => createSomeObservable(count))
            .takeWhile(x => x !== true);
    

    【讨论】:

    • 感谢您的提示。 .take() 在这里工作 (+1)。我用这个作为一个简单的例子;我更大的问题是:如何(如果可能)根据某些条件从第二个可观察到的第一个完成。例如,假设我在一个链中有 2 个 observable,如果第一个 Observable 发出 true,我想停止所有工作。也许takeUntil() 会做什么
    • 2 observables in a chain 是什么意思?如您所料,takeXxx 运营商可能会帮助您。顺便说一句,这篇文章让我对takeXxx 运营商大开眼界。 medium.com/@benlesh/…
    • 我的 OP 在一个链中有 2 个 observables(switchMap 输出第二个)。我只是在寻找一种从其中杀死/完成流的方法,我敢打赌运营商会让我这样做。我会阅读您链接的文章。
    • 我添加了一个示例,如果内部 Observable 发出 true,则终止流。
    • 你是对的,那篇文章以及你在这里的回复给了我我所需要的一切。非常感谢。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-11-19
    • 2022-07-15
    • 2019-01-23
    • 1970-01-01
    • 1970-01-01
    • 2016-11-27
    • 1970-01-01
    相关资源
    最近更新 更多