【问题标题】:Reactive Extensions Buffer on count, interval and event计数、间隔和事件的响应式扩展缓冲区
【发布时间】:2017-10-15 20:37:06
【问题描述】:

我想缓冲发送到我的服务器的事件。刷新缓冲区的触发器是达到缓冲区大小、已达到缓冲期或已卸载窗口。

我通过创建主题并将buffer 与关闭通知程序一起使用来缓冲发送到我的服务器的事件。我使用race 作为关闭通知程序,并使用window.beforeunload 事件来竞争缓冲期。

this.event$ = new Subject();
this.bufferedEvent$ = this.event$
    .buffer(
        Observable.race(
            Observable.interval(bufferPeriodMs),
            Observable.fromEvent(window, 'beforeunload')
        )
    )
    .filter(events => events.length > 0)
    .switchMap(events =>
        ajax.post(
            this.baseUrl + RESOURCE_URL,
            {
                entries: events,
            },
            {
                'Content-Type': 'application/json',
            }
       )
    );

问题是,我现在如何限制缓冲区的大小。即,我从不希望缓冲区有 10 个项目时被刷新。

【问题讨论】:

    标签: javascript rxjs rxjs5


    【解决方案1】:

    使用独立触发器的解决方案让我困扰的一件事是fullBufferTrigger 不知道timeoutTrigger 何时发出它的一个缓冲值,所以给定正确的事件序列,fullBuffer 将在超时后提前触发。

    理想情况下,希望fullBufferTriggertimeoutTrigger 触发时重置,但事实证明这样做很棘手。

    使用bufferTime()

    在 RxJS v4 中有一个运算符 bufferWithTimeOrCount(timeSpan, count, [scheduler]),在 RxJS v5 中它被卷成一个附加签名 bufferTime()(从清晰的角度来看可能是一个错误)。

    bufferTime<T>(
      bufferTimeSpan: number, 
      bufferCreationInterval: number, 
      maxBufferSize: number, 
      scheduler?: IScheduler
    ): OperatorFunction<T, T[]>;
    

    剩下的唯一问题是如何合并window.beforeunload 触发器。查看bufferTime 的源代码,它应该在接收到onComplete 时刷新它的缓冲区。
    因此,我们可以通过向缓冲的事件流发送 onComplete 来处理window.beforeunload

    bufferTime 的规范没有对 onComplete 进行明确的测试,但我想我已经设法将它们放在一起。

    注意事项:

    • 将超时设置为较大,以便将其从图片中取出以进行测试。
    • 源事件流不受影响,说明添加了 event8 但从不发出,因为窗口在它发生之前就被破坏了。
    • 要查看没有 beforeunloadTrigger 的输出流,请注释掉发出onComplete 的行。 Event7 在缓冲区中,但不会发出。

    测试:

    const bufferPeriodMs = 7000  // Set high for this test
    const bufferSize = 2
    const event$ = new Rx.Subject()
    
    /*
      Create bufferedEvent$
    */
    const bufferedEvent$ = event$
      .bufferTime(bufferPeriodMs, null, bufferSize)
      .filter(events => events.length > 0)
    const subscription = bufferedEvent$.subscribe(console.log)  
    
    /*
      Simulate window destroy
    */
    const destroy = setTimeout( () => {
      subscription.unsubscribe()
    }, 4500)
    
    /*
      Simulate Observable.fromEvent(window, 'beforeunload')
    */
    const beforeunloadTrigger = new Rx.Subject()
    // Comment out the following line, observe that event7 does not emit
    beforeunloadTrigger.subscribe(x=> event$.complete())
    setTimeout( () => {
      beforeunloadTrigger.next('unload')
    }, 4400)
    
    /*
      Test sequence
      Event stream:        '(123)---(45)---6---7-----8--|'
      Destroy window:      '-----------------------x'
      window.beforeunload: '---------------------y'
      Buffered output:     '(12)---(34)---(56)---7'
    */
    event$.next('event1')
    event$.next('event2')
    event$.next('event3')
    setTimeout( () => { event$.next('event4'); event$.next('event5') }, 1000)
    setTimeout( () => { event$.next('event6') }, 3000)
    setTimeout( () => { event$.next('event7') }, 4000)
    setTimeout( () => { event$.next('event8') }, 5000)
    

    工作示例:CodePen

    【讨论】:

    • 谢谢!他们的医生真的不清楚!这么浪费时间
    【解决方案2】:

    这是我的工作解决方案。添加了额外的 console.log() 以显示事件的顺序。

    唯一有点麻烦的是fullBufferTrigger中的.skip(1),但它是必需的,因为它会在缓冲区满时触发(natch),但bufferedEvent$中的缓冲区之前似乎没有最新事件它被触发了。

    幸运的是,有了timeoutTrigger,最后一个事件就会发出。如果没有超时,fullBufferTrigger 本身不会发出最终事件。

    另外,将buffer 更改为bufferWhen,因为前者似乎没有通过两个触发器触发,尽管您希望从文档中得到它。
    脚注buffer(race()) 的比赛只完成一次,所以无论哪个触发器先到达那里都会被使用,而其他触发器将被忽略。相比之下,bufferWhen(x =&gt; race()) 会在每次事件发生时进行评估。

    const bufferPeriodMs = 1000
    
    const event$ = new Subject()
    event$.subscribe(event => console.log('event$ emit', event))
    
    // Define triggers here for testing individually
    const beforeunloadTrigger = Observable.fromEvent(window, 'beforeunload')
    const fullBufferTrigger = event$.skip(1).bufferCount(2)
    const timeoutTrigger = Observable.interval(bufferPeriodMs).take(10)
    
    const bufferedEvent$ = event$
      .bufferWhen( x => 
        Observable.race(
          fullBufferTrigger,
          timeoutTrigger
        )
      )
      .filter(events => events.length > 0)
    
    // output
    fullBufferTrigger.subscribe(x => console.log('fullBufferTrigger', x))
    timeoutTrigger.subscribe(x => console.log('timeoutTrigger', x))
    bufferedEvent$.subscribe(events => {
      console.log('subscription', events)
    })
    
    // Test sequence
    const delayBy = n => (bufferPeriodMs * n) + 500 
    event$.next('event1')
    event$.next('event2')
    event$.next('event3')
    
    setTimeout( () => {
      event$.next('event4')
    }, delayBy(1))
    
    setTimeout( () => {
      event$.next('event5')
    }, delayBy(2))
    
    setTimeout( () => {
      event$.next('event6')
      event$.next('event7')
    }, delayBy(3))
    

    工作示例:CodePen

    编辑:触发缓冲区的替代方法

    由于bufferWhenrace 的组合可能有点低效(每次事件发射都会重新开始比赛),另一种方法是将触发器合并到一个流中并使用简单的buffer

    const bufferTrigger$ = timeoutTrigger
      .merge(fullBufferTrigger)
      .merge(beforeunloadTrigger)
    
    const bufferedEvent$ = event$
      .buffer(bufferTrigger$)
      .filter(events => events.length > 0)
    

    【讨论】:

    • 我知道你的目标是什么。我尝试了类似的方法。您需要先进行计数。像event$.count().filter(c =&gt; c &gt; 2) 这样的东西。但问题是 Observable 不会产生任何东西,因为没有订阅。
    • 那么,您正在使用event$.next(event) 参与活动?
    • event$.bufferCount(10) 怎么样?这需要测试。
    • 是的,我正在使用 event$.next(event)。我尝试使用 bufferCount ......它有点工作,问题是当任何其他事件发生时它不会重置为 0
    • 我认为在这种情况下使用 race 存在根本缺陷,因为它选择了第一个发出并坚持它的 observable。 bufferWhen 通过调用每个事件的函数来解决问题,因此重复每个事件的比赛(效率低吗?)。我实际上是在查看buffer(fullBufferTrigger.merge(timeoutTrigger).merge(beforeunloadTrigger)),这更符合您的目标。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-15
    • 2013-07-22
    相关资源
    最近更新 更多