【问题标题】:how to buffer the latest value until another value arrive in another sequence by rxjs?如何缓冲最新值,直到另一个值通过rxjs以另一个序列到达?
【发布时间】:2017-10-14 02:10:16
【问题描述】:

我正在尝试在我的项目中使用 rxjs。我有以下序列,我期望的是第一个序列只会在一个值到达另一个序列后处理,并且只保留第一个序列中的最新值。有什么建议吗?

s1$ |a---b----c-

s2$ |------o----

预期结果:

s3$ |------b--c-

【问题讨论】:

  • 看起来你想要takeUntil() 或skipUntil,但我想我不明白你的描述。您无法获得预期的------|b-c-,因为完整的信号始终是最后一个信号。换句话说,您不能在发送完整信号后发出更多值。
  • @martin 我的描述造成了一些混乱,对此感到抱歉。那个|并不意味着完整的信号,我更新了我的描述。我是 rxjs 的新手。如果我使用takeUtil,我想我会失去b,并且s1.subcribe只会在c到达后触发,是吗?

标签: javascript typescript rxjs5


【解决方案1】:

我想我会使用 ReplaySubject 来做到这一点

const subject$ = new Rx.ReplaySubject(1)

const one$ = Rx.Observable.interval(1000) 
const two$ = Rx.Observable.interval(2500)

one$.subscribe(subject$)

const three$ = two$
  .take(1)
  .flatMap(() => subject$)

// one$   |----0---1---2---3---4---
// two$   |----------0---------1---
// three$ |----------1-2---3---4---

【讨论】:

    【解决方案2】:

    我会将已经非常相似的sample() 和skipUntil() 结合起来。

    const start = Scheduler.async.now();
    const trigger = new Subject();
    
    const source = Observable
        .timer(0, 1000)
        .share();
    
    Observable.merge(source.sample(trigger).take(1), source.skipUntil(trigger))
        .subscribe(val => console.log(Scheduler.async.now() - start, val));
    
    setTimeout(() => {
        trigger.next();
    }, 2500);
    

    这将输出以2 开头的数字。

    source  0-----1-----2-----3-----4
    trigger ---------------X---------
    output  ---------------2--3-----4
    

    带有时间戳的控制台输出:

    2535 2
    3021 3
    4024 4
    5028 5
    

    您也可以使用switchMap() 和ReplaySubject,但它可能不如前面的示例那么明显,您需要两个主题。

    const start = Scheduler.async.now();
    const trigger = new Subject();
    
    const source = Observable
        .timer(0, 1000)
        .share();
    
    const replayedSubject = new ReplaySubject(1);
    source.subscribe(replayedSubject);
    
    trigger
        .switchMap(() => replayedSubject)
        .subscribe(val => console.log(Scheduler.async.now() - start, val));
    
    setTimeout(() => {
        trigger.next();
    }, 2500);
    

    输出完全一样。

    【讨论】:

      【解决方案3】:

      last + takeUntil 会起作用

      这是一个例子:

      let one$ = Rx.Observable.interval(1000);
      let two$ = Rx.Observable.timer(5000, 1000).mapTo('stop');
      
      one$
        .takeUntil(two$)
        .last()
        .subscribe(
           x=>console.log(x),
           err =>console.error(err),
           ()=>console.log('done')
        );
      

      【讨论】:

      • 您的视频流结束并且没有采用最新的值。
      猜你喜欢
      • 1970-01-01
      • 2020-09-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-09-15
      • 1970-01-01
      • 1970-01-01
      • 2019-06-15
      相关资源
      最近更新 更多