【问题标题】:How can I make an rxjs Subject repeat its last emission periodically?如何让 rxjs 主题定期重复其最后一次发射?
【发布时间】:2017-04-18 02:44:17
【问题描述】:

我正在尝试创建一个 rxjs 主题,该主题将在一段时间不活动后重复其输出。我最初的设计是利用 debounceTime,但是这似乎不会多次触发。

我希望主题在调用next 时立即发射,并定期重复发射,直到提供新值:

Inputs:  ---a---------b---------c----
Outputs: ---a---a---a-b---b---b-c---c

目前我有这样的事情:

const subject = new rx.Subject()
subject.debounceTime(5000)
       .subscribe(subject)

subject.subscribe(value => console.log(`emitted: ${value}`))
subject.take(1).subscribe(next => next, error => error, () => {
    console.log('emitted once')
})
subject.take(2).subscribe(next => next, error => error, () => {
    console.log('emitted twice')
})
subject.take(3).subscribe(next => next, error => error, () => {
    console.log('emitted thrice')
})

subject.next('a')

但是,这只会发出一次“a”,并且永远不会看到输出“发出三次”。

有人可以帮我理解这里出了什么问题吗?

【问题讨论】:

    标签: javascript rxjs


    【解决方案1】:

    如果我正确理解您的问题,我认为您可以使用 repeatWhen() 运算符:

    const subject = new ReplaySubject(1);
    subject.next('a');
    
    subject
      .take(1)
      .repeatWhen(() => Observable.timer(500, 500))
      .subscribe(val => console.log(val));
    
    setTimeout(() => subject.next('b'), 1400);
    

    观看现场演示:https://jsbin.com/nufasiq/2/edit?js,console

    这会以 500 毫秒的间隔打印以控制台以下输出:

    a
    a
    a
    b
    b
    b
    

    take(1) 在这里是必要的,以使链正确完成,它被再次订阅其源 Observable 的repeatWhen() 拦截。

    【讨论】:

    • 这太好了,谢谢。 repeatWhen() 似乎确实解决了一半的问题,但是它似乎在计时器发出时发出。我已经合并了原始主题并消除了中继器的抖动,以产生足够接近我需要的东西。 subject.merge(subject.debounceTime(1500).take(1).repeatWhen(() => Observable.timer(1500, 1500)))
    【解决方案2】:

    另一种选择是使用switchMapinterval

    const source = new Rx.Subject();
    
    source
      .switchMap((val) => Rx.Observable.interval(5000).map(() => val))
      .subscribe(val => console.log(val));
    
    
    Rx.Observable.interval(11000).subscribe(x => source.next(x));
    

    查看demo in jsbin

    【讨论】:

      【解决方案3】:

      你可以尝试类似this jsfiddle(rxjs v4,但我想用switchMap替换flatMapLatest应该对v5有用):

      const input$ = new Rx.Subject()
      
      const repeatedInput$ = input$.flatMapLatest(input => Rx.Observable.interval(200).map(input))
      output$ = Rx.Observable.merge(input$, repeatedInput$)
      
      output$.subscribe(console.log.bind(console, `output`))
      input$.onNext(2)
      setTimeout(() => input$.onNext(4), 300)
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-11-22
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-10-29
        • 1970-01-01
        • 2019-05-23
        • 1970-01-01
        相关资源
        最近更新 更多