【问题标题】:Resubscribing same shared observable inside flatMap of own data emits no data. By design?在自己数据的 flatMap 中重新订阅相同的共享 observable 不会发出任何数据。按设计?
【发布时间】:2017-09-06 14:13:37
【问题描述】:

尝试迁移到 rx-java2 并在重新订阅它自己的 flatMap 内的共享 observable 时遇到了问题。需要这种模式来获取更新刷新链:

  1. 从网络获取当前数据(共享 observable 以避免多个网络请求,如果源同时被多个观察者订阅)。
  2. 修改数据并发回服务器(可完成)
  3. 更新完成后再次获取数据

整个过程是这样的:

@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) 从顶部重新启动 observable o
  • @Bob 选项 c。应该是一个刷新的值并且之前工作过(由于我猜的实现细节)。我有一个基于连接各个可观察对象的网络磁盘内存模型,而不是真实应用程序中的o。网络一是共享的,以防止多个网络请求。数据修改请求后,flatmap 内的缓存正在清除。然后,如果从网络读取原始数据,则第二个订阅将获得相同的共享源,因为 take(1)flatMap 正文评估时不会取消订阅原始数据,而后一个订阅仅获得 onComplete
  • 发出单个值的共享 observable 是我的代码中的一个设计缺陷,因为它容易出现竞争条件。订阅者可能介于价值发射和完全发射之间并得到一个空结果。在现实世界中没有看到它,但上面的测试清楚地表明了这一点。

标签: rx-java2


【解决方案1】:

看看“分享”的气泡图,您就会明白为什么它会这样:Observable.share()

share() 发出订阅后发出的项目,它不会重新发出以前发出的项目。查看Observable.replay(),了解您所期望的行为。

【讨论】:

  • 没有。重播不是预期的 - 在flatMap 内需要刷新数据。但我不知道如何取消订阅源以停止共享原始 observable(重新运行副作用)。我添加了take(1) 试图澄清。似乎 take(1) 在 2.x 之前急切地取消订阅源,并且从现在开始一直在 flatMap 点等待......
  • 我的意思是我不(不想)知道我身边的源是热的还是冷的。并且无法弄清楚如何确保(使其)冷并重做副作用
【解决方案2】:

似乎不是“重新启动”共享可观察对象及其副作用的合法方式,因为不能保证其他订阅者目前没有在收听。

【讨论】:

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