【问题标题】:RxJava buffering - ignoring zero itemsRxJava 缓冲 - 忽略零项
【发布时间】: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 -&gt; !collection.isEmpty)之后使用.distinctUntilChanged吗?

标签: java rx-java reactive-programming reactivex


【解决方案1】:

正如 Mark Keen 所建议的,.distinctUntilChanged 可以解决问题。

所以如果缓冲后有1+项,下面的代码会推送事件列表:

connectableObservable.toFlowable(BackpressureStrategy.BUFFER)
    .filter(event -> event.getId().equals(eventId))
    .buffer(1, TimeUnit.SECONDS, 50)
    .distinctUntilChanged()             // <<<======  
    .map(eventsMapper::mapCollection);

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-10-10
    • 2014-02-25
    • 2019-04-01
    • 1970-01-01
    • 2014-04-23
    • 1970-01-01
    相关资源
    最近更新 更多