【问题标题】:Rxjs: Observable with takeUntil(timer) keeps emitting after the timer has tickedRxjs:带有takeUntil(timer)的Observable在计时器计时后继续发射
【发布时间】:2017-01-09 07:36:38
【问题描述】:

我遇到了takeUntil() 的一个非常奇怪的行为。我创建了一个可观察的计时器:

let finish = Observable.timer(3000);

那我等一会儿再打电话

// 2500 ms later
someObservable.takeUntil(finish);

我希望所说的 observable 在计时器“滴答”之后停止发射,即在它创建后大约 500 毫秒。实际上它在创建后持续发射 3000 毫秒,远远超过计时器“滴答”的那一刻。如果我使用包含绝对时间值的 Date 对象创建计时器,则不会发生这种情况。

这是设计使然吗?如果有,解释是什么?

这里是完整的代码,可以用 node.js 运行(它需要npm install rx):

let {Observable, Subject} = require("rx")
let start = new Date().getTime();
function timeMs() { return new Date().getTime() - start };

function log(name, value) { 
    console.log(timeMs(), name, value);
}

Observable.prototype.log = function(name) {
    this.subscribe( v=>log(name,v), 
                    err=>log(name, "ERROR "+err.message), 
                    ()=>log(name, "DONE"));
    return this;
}

let finish = Observable.timer(3000).log("FINISH");
setTimeout( ()=>Observable.timer(0,500).takeUntil(finish).log("seq"), 2500);

这会生成以下输出:

2539 'seq' 0
3001 'FINISH' 0
3005 'FINISH' 'DONE'
3007 'seq' 1
3506 'seq' 2
4006 'seq' 3
4505 'seq' 4
5005 'seq' 5
5506 'seq' 6
5507 'seq' 'DONE'

如果我使用绝对时间创建计时器:

let finish = Observable.timer(new Date(Date.now()+3000)).log("FINISH");

然后它按预期运行:

2533 'seq' 0
3000 'seq' 'DONE'
3005 'FINISH' 0
3005 'FINISH' 'DONE'

这种行为在各种情况下似乎相当一致。例如如果您使用mergeMap()switchMap() 间隔并创建子序列,结果将是相似的:子序列在完成事件之后继续发射。

想法?

【问题讨论】:

    标签: javascript rxjs


    【解决方案1】:

    你忘记了Observables的第一条冷规则:每个订阅都是一个新流。

    您的log 运算符有错误;它订阅了一次Observable(从而创建了第一个订阅),然后返回原始Observable,当您将其传递给takeUntil 运算符时,它会隐式地再次订阅Observable .因此,实际上您实际上有两个活动的seq 流,它们的行为都是正确的。

    它适用于绝对情况,因为您基本上将每个流设置为在特定时间发出,而不是订阅发生时的相对时间。

    如果您想看到它的工作,我建议您将实现更改为:

    let start = new Date().getTime();
    function timeMs() { return new Date().getTime() - start };
    
    function log(name, value) { 
        console.log(timeMs(), name, value);
    }
    
    Observable.prototype.log = function(name) {
        // Use do instead of subscribe since this continues the chain
        // without directly subscribing.
        return this.do(
          v=>log(name,v), 
          err=>log(name, "ERROR "+err.message), 
          ()=>log(name, "DONE")
        );
    }
    
    let finish = Observable.timer(3000).log("FINISH");
    
    setTimeout(()=> 
      Observable.timer(0,500)
        .takeUntil(finish)
        .log("seq")
        .subscribe(), 
    2500);
    

    【讨论】:

    • 你给了我太多的信任:我一开始就不知道冷可观察的规则:) 你是说 timer() 只有在订阅后才开始计算时间,并且计数是否为每个新订阅重新开始?坦率地说,timer() 方法的文档中没有任何内容暗示这一点,除了可能提到“async IScheduler”,它可能是从 .NET 实现中复制的,在 JavaScript 上下文中对我来说毫无意义。
    • 别担心,这是一个很常见的错误。它不是 timer 运算符特有的东西,它是 Observables 固有的更一般的行为。它们旨在被懒惰地评估。在hot vs. cold 上看到这个相当不错的入门
    【解决方案2】:

    为了完整起见,这里是实际执行我想要的代码。通过使用Observable.publish().connect(),它创建了一个“热”计时器,该计时器立即开始计时,并为所有订阅者保持相同的时间。正如@paulpdaniels 所建议的,它还避免了“日志”方法中不需要的订阅。

    警告:注意竞态条件。如果子序列在计时器计时后开始,它将永远不会停止。为了演示,将最后一行的超时时间从 2500 更改为 3500。

    let {Observable, Subject, Scheduler, Observer} = require("rx")
    let start = new Date().getTime();
    function timeMs() { return new Date().getTime() - start };
    
    function log(name, value) { 
        console.log(timeMs(), name, value);
    }
    
    var logObserver =  function(name) {
        return Observer.create( 
          v=>log(name,v), 
          err=>log(name, "ERROR "+err.message), 
          ()=>log(name, "DONE"));
    }
    
    Observable.prototype.log = function(name) { return this.do(logObserver(name)); }
    
    Observable.prototype.start = function() { 
        var hot = this.publish(); hot.connect(); 
        return hot; 
    }
    
    let finish = Observable.timer(3000).log("FINISH").start();
    
    setTimeout(()=> 
      Observable.timer(0,500)
        .takeUntil(finish)
        .log("seq")
        .subscribe(), 
    2500);
    

    输出是

    2549 'seq' 0
    3002 'FINISH' 0
    3006 'seq' 'DONE'
    3011 'FINISH' 'DONE'
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-08-31
      • 1970-01-01
      • 2017-07-11
      • 1970-01-01
      • 2012-07-18
      • 2017-11-06
      • 1970-01-01
      相关资源
      最近更新 更多