【问题标题】:How to remember state with retry operators in RxJava2如何在 RxJava2 中使用重试运算符记住状态
【发布时间】:2017-11-02 03:01:53
【问题描述】:

我有一个能够从中断中恢复的网络客户端,但在重试时需要最后一条消息。

Kotlin 中的示例:

fun requestOrResume(last: Message? = null): Flowable<Message> =
    Flowable.create({ emitter ->
        val connection = if (last != null)
                             client.start()
                         else
                             client.resumeFrom(last.id)

        while (!emitter.isDisposed) {
            val msg = connection.nextMessage()
            emitter.onNext(msg)
        }
    }, BackpressureStrategy.MISSING)

requestOrResume()
    .retryWhen { it.flatMap { Flowable.timer(5, SECONDS) } }
    // how to pass the resume data when there is a retry?

问题:如您所见,我需要最后收到的消息来准备恢复电话。如何跟踪它,以便在重试时可以提出恢复请求?

一种可能的解决方案可能是创建一个持有者类,它只保存对最后一条消息的引用,并在收到新消息时更新。这样,当重试时,可以从持有者那里获得最后一条消息。示例:

class MsgHolder(var last: Message? = null)

fun request(): Flowable<Message> {
    val holder = MsgHolder()
    return Flowable.create({ emitter ->
        val connection = if (holder.last != null)
                             client.start()
                         else
                             client.resumeFrom(holder.last.id)

        while (!emitter.isDisposed) {
            val msg = connection.nextMessage()
            holder.last = msg // <-- update holder reference
            emitter.onNext(msg)
        }
    }, BackpressureStrategy.MISSING)
}

我认为这可能有效,但感觉就像是 hack(线程同步问题?)。

是否有更好的方法来跟踪状态以便重试?

【问题讨论】:

    标签: android kotlin rx-java rx-java2


    【解决方案1】:

    请注意,除非您在最后一个元素周围重新抛出一个包装器(在功能上与您现有的“hack”-ish 解决方案没有太大不同,但在 imo 上更丑),否则没有任何错误处理操作员可以在没有外部帮助的情况下恢复最后一个元素,因为它们只能访问Throwable 的流。相反,请查看以下递归方法是否适合您的需求:

    fun retryWithLast(seed: Flowable<Message>): Flowable<Message> {
      val last$ = seed.last().cache();
      return seed.onErrorResumeNext {
        it.flatMap {
          retryWithLast(last$.flatMap {
            requestOrResume(it)
          })
        }
      };
    }
    retryWithLast(requestOrResume());
    

    最大的区别是使用cache 将上次尝试的尾随值缓存在 Observable 中,而不是手动在值中这样做。另请注意,错误处理程序中的递归意味着如果后续尝试继续失败,retryWithLast 将继续扩展流。

    【讨论】:

      【解决方案2】:

      仔细查看buffer()运营商:link 你可以这样使用它:

      requestOrResume()
          .buffer(2)
      

      从现在开始,您的 Flowable 将返回带有两个最新对象的 List&lt;Message&gt;

      【讨论】:

      • 我看不出buffer 运算符在这种情况下有何帮助。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-03-28
      相关资源
      最近更新 更多