【问题标题】:RxJava bufferWithTimeOrCount() implementation?RxJava bufferWithTimeOrCount() 实现?
【发布时间】:2017-01-05 17:29:13
【问题描述】:

使用 RxJava,我需要将一个项目流缓冲成 3 个一组,但如果传入项目之间的时间超过 500 毫秒,则刷新缓冲区。

bufferWithTimeOrCount() 运算符正是我正在寻找的,但它似乎只为RxJSRx.NET 实现,我需要为此使用RxJava。

有没有办法复制 bufferWithTimeOrCount() 的行为并使用现有的 RxJava 1.x 运算符获得我想要的?

buffer(500, TimeUnit.MILLISECONDS, 3) 尝试每 500 毫秒发出一个新列表,不考虑自上一项以来的时间。

我一直在努力使事情与buffer(bufferClosingSelector) 一起工作(请参阅RxJava Buffer)。我设置了一个序列,它发出项目之间的延迟增加:

Observable<Long> emitter = Observable
    .range(1, 9)
    .flatMap(n -> {
        return Observable.just(n).delay(n * n * 50, TimeUnit.MILLISECONDS); 
    })
    .timeInterval()
    .map(interval -> interval.getIntervalInMilliseconds());

然后我尝试在项目之间的时间超过 500 毫秒时使用 debounce() 刷新缓冲区:

emitter
    .buffer(emitter.debounce(500, TimeUnit.MILLISECONDS))
    .toBlocking()
    .subscribe(i -> System.out.println(i));

这似乎可行,产生如下序列:

[51, 150, 250, 352, 449]
[551]
[650]
[749]
[850]

然后我尝试创建一个计数器来在每三个项目之后刷新缓冲区,认为我可以使用去抖动的 observable merge() 它:

emitter
    .buffer(emitter
       .scan(1L, (n, x) -> n + 1)   // count the items up from 1
       .filter(n -> ((n % 3) == 0)  // emit every three items
    )
    .toBlocking()
    .subscribe(i -> System.out.println(i));

但结果并不那么成功:

[50]
[152, 248, 350]
[452, 550, 650, 749]
[]

如果我输出间隔,则计数器似乎与批处理 observable 异步运行,因此在它触发时间和批处理实际释放时间之间会发生竞争条件:

emitter
    .buffer(emitter
       .doOnNext(i -> System.out.println("i = " + i))
       .scan(1L, (n, x) -> n + 1)   // count the items up from 1
       .filter(n -> ((n % 3) == 0)  // emit every three items
    )
    .toBlocking()
    .subscribe(i -> System.out.println(i));

生产:

i = 63
i = 150
[51, 150]
i = 249
i = 350
i = 451
[249, 351]
i = 550
i = 650
i = 750
[450, 550, 650, 750]
i = 850
[850]

...现在我有点难过。

这很长,但是 TL;DR:有没有办法复制 bufferWithTimeOrCount() 的行为并使用现有的 RxJava 1.x 运算符获得我所追求的?谢谢!

【问题讨论】:

  • 在 2.x 中,有一个重载版本可以让你重新启动时间窗口:buffer(long timespan, TimeUnit unit, Scheduler scheduler, int count, Callable&lt;U&gt; bufferSupplier, boolean restartTimerOnMaxSize)

标签: rx-java


【解决方案1】:

编辑:再次查看问题后,我意识到这个解决方案总是会为每个第三个项目发射,无论它是否已经由于去抖动而发射。我无法立即想到在去抖动后“重置”计数缓冲区的方法。

问题在于,通过使用 emitter 两次,您会创建两个具有不确定行为的独立流,因为涉及到调度程序。

关键是结合buffer(...)debounce(...)publish(...)运算符:

emitter.publish(source -> {
    // Create the buffer here using a shared source
    return source.buffer(Observable.defer(() ->
         // Merge these reasons for closing the buffer
         Observable.merge(
             // Either after 500 ms
             source.debounce(500, TimeUnit.MILLISECONDS),
             // or every 3 items
             source.buffer(3)
         )
    ))
})

我无法测试此解决方案,但 cmets 应该解释其背后的想法。 Observable.defer 确保负责关闭窗口的可观察对象是在原始源之后创建的。

【讨论】:

  • 确实如此。使用scan() 方法对项目进行计数(就像我在问题中提到的那样)打开了附加takeUntil(debounceObservable).repeat() 以在发生去抖动时重置计数器的可能性。但这首先需要计数器可靠地工作。
猜你喜欢
  • 1970-01-01
  • 2023-03-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多