【问题标题】:Accumulate and Emit behaviour in RxJava 2RxJava 2 中的 Accumulate 和 Emit 行为
【发布时间】:2019-06-21 16:12:18
【问题描述】:

我有一个Status 类代表一些任务状态。我有一个带有以下签名的方法

def update(existing: Status, update: Status): Status

该函数接受旧状态和新更新,并返回结合这两者的新Status 对象。

我希望自定义主题的行为如下:

  • 写行为

    • 它应该包含一个Status 的实例。我们就叫它current
    • 我应该能够订阅给定的Flowable<Status>
    • 在每次新的状态更新(通过上述Flowable 提供)时,它都会调用上述update 方法并将结果存储在current 中。
  • 阅读行为

    • 定期读取,应该能够定期发出当前状态。类似Flowable.interval(1 second).map(i->current)

我可以使用StatusHolder 类来实现上述目的,用于写入行为(它负责保存和更新)以及单独订阅Flowable.interval(1 second).map(i->holder.current)

在实现这个之后,我遇到了Subject的概念,它既是Observable又是Observer。我想要的功能与此类似。我想要一个可以接收和发出Status 对象但需要在接收时进行一些计算的类。

我查看了现有的Subject 实现,我认为它们中的任何一个都不自然地支持这种行为。第二件事是,它们在Observable 上运行,而不是Flowable,所以我需要使用toFlowabletoObservable 来使用它。

有没有更好的方法来实现这种行为?

【问题讨论】:

    标签: rx-java reactive-programming rx-java2 reactive


    【解决方案1】:

    让我们将source 定义为您的来源Flowable<Status>

    Flowable<Status> source = ...
    

    那你可以试试这个:

    Flowable.interval(1, TimeUnit.SECONDS)
            .withLatestFrom(source.scan(this::update), (i, status) -> status)
            .share();
    
    • .interval() 运算符会定期发出一个 long 值。
    • .scan() 运算符“累积”源流发出的项目。在您的情况下,累加器是 .update 函数
    • .withLatestFrom() 结合了两个流。这将替换您的 .map(i-&gt;holder.current)
    • .share() 允许所有订阅者共享一个订阅。


    附注:

    Subject 有一个Flowable 版本,称为Processor

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-08-17
      • 2018-11-04
      • 1970-01-01
      • 2021-11-30
      • 1970-01-01
      相关资源
      最近更新 更多