【发布时间】:2019-01-19 21:01:54
【问题描述】:
这是我用于缓冲和转换传入事件的代码:
public Publisher<Collection<EventTO>> logs(String eventId) {
ConnectableObservable<Event> connectableObservable = eventsObservable
.share().publish();
connectableObservable.connect();
connectableObservable.toFlowable(BackpressureStrategy.BUFFER)
.filter(event -> event.getId().equals(eventId))
.buffer(1, TimeUnit.SECONDS, 50)
.map(eventsMapper::mapCollection);
}
这里的问题是Flowable 每秒返回一个空列表,尽管没有事件发布到eventsObservable。
有没有办法保持.buffer 直到至少有一个对象?
注意: 看起来有一种方法可以在 C# 中做到这一点(在此处描述:https://stackoverflow.com/a/30090185/668148)。 但是如何用 Java 来实现呢?
【问题讨论】:
-
你能在
buffer( ... )或.filter(collection -> !collection.isEmpty)之后使用.distinctUntilChanged吗?
标签: java rx-java reactive-programming reactivex