【发布时间】:2017-09-06 14:13:37
【问题描述】:
尝试迁移到 rx-java2 并在重新订阅它自己的 flatMap 内的共享 observable 时遇到了问题。需要这种模式来获取更新刷新链:
- 从网络获取当前数据(共享 observable 以避免多个网络请求,如果源同时被多个观察者订阅)。
- 修改数据并发回服务器(可完成)
- 更新完成后再次获取数据
整个过程是这样的:
@Test fun sharedTest() {
val o = Observable.just(1).share()
assertEquals(1, o
.take(1)
.flatMap({
Completable.complete()
.andThen(o) })
.blockingFirst())
}
测试失败:java.util.NoSuchElementException
如果o 未共享,则一切正常。
这种行为似乎是因为当原始的单个值已经被调度并且只有onComplete 事件被看到时,后一个订阅者来了。
有人知道这是一种设计行为并以某种方式记录吗?当然有一个解决方法,但我需要知道原因,因为这有点烦人。该方法在 Rx 1.x 中有效
目前使用的是 2.1.3 版
编辑:
似乎没有合法的方式来“重启”一个共享的 observable 及其副作用,因为不能保证其他订阅者目前没有在听。
【问题讨论】:
-
我不清楚
flatMap()的输出是否应该是a)从o获取的原始值,或者b)o的新值已修改,或者 c) 从顶部重新启动 observableo。 -
@Bob 选项 c。应该是一个刷新的值并且之前工作过(由于我猜的实现细节)。我有一个基于连接各个可观察对象的网络磁盘内存模型,而不是真实应用程序中的
o。网络一是共享的,以防止多个网络请求。数据修改请求后,flatmap内的缓存正在清除。然后,如果从网络读取原始数据,则第二个订阅将获得相同的共享源,因为take(1)在flatMap正文评估时不会取消订阅原始数据,而后一个订阅仅获得onComplete。 -
发出单个值的共享 observable 是我的代码中的一个设计缺陷,因为它容易出现竞争条件。订阅者可能介于价值发射和完全发射之间并得到一个空结果。在现实世界中没有看到它,但上面的测试清楚地表明了这一点。
标签: rx-java2