【问题标题】:RxJava buffer/window with custom counting criteria具有自定义计数标准的 RxJava 缓冲区/窗口
【发布时间】:2018-10-13 09:53:09
【问题描述】:

我有一个 Observable,它发出许多对象,我想使用 windowbuffer 操作对这些对象进行分组。但是,我希望能够使用自定义条件,而不是指定 count 参数来确定窗口中应该有多少对象。

例如,假设 observable 发出 Message 类的实例,如下所示。

class Message(
   val int size: Int
)

我想根据它们的size 变量来缓冲或窗口化消息实例,而不仅仅是它们的计数。例如,要获得总大小最多为 5000 的消息窗口。

// Something like this
readMessages()
    .buffer({ message -> message.size }, 5000)

有没有简单的方法可以做到这一点?

【问题讨论】:

标签: kotlin rx-java rx-java2


【解决方案1】:

首先我必须承认,我不是 RxJava 专家。 我刚刚发现您的问题具有挑战性,并试图找到解决方案。

有一个带有参数boundaryIndicatorwindow() 函数。如果达到窗口大小,您必须创建一个发出项目的Publisher/Flowable

在示例中,我创建了一个对象windowManager,用作boundaryIndicator。在onNext 回调中,我调用windowManager 并让它有机会打开一个新窗口。

val windowManager = object {
    lateinit var emitter: FlowableEmitter<Unit>
    var windowSize: Long = 0

    fun createEmitter(emitter: FlowableEmitter<Unit>) {
        this.emitter = emitter
    }

    fun openWindowIfRequired(size: Long) {
        windowSize += size
        if (windowSize > 5) {
            windowSize = 0
            emitter.onNext(Unit)
        }
    }
}

val windowBoundary = Flowable.create<Unit>(windowManager::createEmitter, BackpressureStrategy.ERROR)

Flowable.interval(1, TimeUnit.SECONDS).window(windowBoundary).subscribe {
    it.doOnNext {
        windowManager.openWindowIfRequired(it)
    }.doOnSubscribe {
        println("Open window")
    }.doOnComplete {
        println("Close window")
    }.subscribe {
        println(it)
    }
}

【讨论】:

  • 谢谢!这似乎是一个非常常见的场景(在我的情况下,我正在聚合传入的蓝牙数据包),但文档并没有提供太多关于如何处理这个问题的指导。也许把它变成一个专用的 Rx 方法是有意义的,……不知何故。作为旁注,我不得不依赖doAfterNext,因为我希望将触发新窗口的发射项目包含在窗口中。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2011-10-03
  • 2018-11-30
  • 1970-01-01
  • 2020-02-26
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多