【发布时间】:2017-01-05 17:29:13
【问题描述】:
使用 RxJava,我需要将一个项目流缓冲成 3 个一组,但如果传入项目之间的时间超过 500 毫秒,则刷新缓冲区。
bufferWithTimeOrCount() 运算符正是我正在寻找的,但它似乎只为RxJS 和Rx.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<U> bufferSupplier, boolean restartTimerOnMaxSize)。
标签: rx-java