【发布时间】:2022-01-07 18:39:05
【问题描述】:
当数据传输到我的应用程序时,它遵循以下顺序:
- 1 个带有 ID 信息的 ReadStart 数据包
- 1 个或多个 DataPackets 组合形成有效负载
- 1 个 ReadDone 数据包表示传输完成
我有一个创建 Observable 的 Kotlin RX 函数:
val readStartPublishProcessor: PublishProcessor<ReadStartPacket>
val dataPacketPublishProcessor: PublishProcessor<DataPacket>
val readDonePublishProcessor: PublishProcessor<ReadDonePacket>
...
private fun setupReadListener(): Flowable<ReadEvent> {
val dataFlowable = dataPacketPublishProcessor.buffer(readDonePublishProcessor)
return readStartPublishProcessor
.zipWith(other = dataFlowable) { readStart, dataPackets ->
Log.d(tag, "ReadEvent with ${dataPackets.size} data packets")
ReadEvent(event = readStart, payload = combinePackets(dataPackets))
}
}
通过阅读.zipWith 的文档,我希望.zipWith 函数体对readStartPublishProcessor 和dataFlowable 发出的每对值执行一次,然后将该计算结果传递给每个订阅者:
Zip 方法返回一个 Observable,它应用你的函数 选择由两个(或 更多)其他 Observables,这个函数的结果变成 返回的 Observable 发出的项目。 ...它只会发出 many items 作为源 Observable 发出的项目数 发出最少的项目。
但如果我有超过 1 个观察者,我会看到 .zipWith 函数体执行的次数与观察者的数量相同,每次都使用相同的一对发射值。这是一个问题,因为从 .zipWith 函数体中调用的函数会产生副作用。 (注意:观察者中不使用 .share 和 .replay 运算符。)
为什么似乎为每个观察者运行 .zipWith 函数体而不是只运行一次,有没有办法编写它以便它只执行一次而不管观察者的数量?
【问题讨论】: