【问题标题】:Split Rx Observable into multiple streams and process individually将 Rx Observable 拆分为多个流并单独处理
【发布时间】:2015-03-04 12:36:26
【问题描述】:

这是我正在尝试完成的图片。

--a-b-c-a--bbb--a

拆分成

--a-----a--------a --> 一个流

----b------bbb--- --> b流

------c---------- --> c流

那么,可以

a.subscribe()
b.subscribe()
c.subscribe()

到目前为止,我发现的所有内容都使用 groupBy() 拆分流,但随后将所有内容折叠回单个流并在同一个函数中处理它们。我想做的是以不同的方式处理每个派生流。

我现在的做法是做一堆过滤器。有没有更好的方法来做到这一点?

【问题讨论】:

    标签: reactive-programming rx-java rxjs


    【解决方案1】:

    简单易行,只需使用filter

    scala 中的一个例子

    import rx.lang.scala.Observable
    
    val o: Observable[String] = Observable.just("a", "b", "c", "a", "b", "b", "b", "a")
    val hotO: Observable[String] = o.share
    val aSource: Observable[String] = hotO.filter(x ⇒ x == "a")
    val bSource: Observable[String] = hotO.filter(x ⇒ x == "b")
    val cSource: Observable[String] = hotO.filter(x ⇒ x == "c")
    
    aSource.subscribe(o ⇒ println("A: " + o), println, () ⇒ println("A Completed"))
    
    bSource.subscribe(o ⇒ println("B: " + o), println, () ⇒ println("B Completed"))
    
    cSource.subscribe(o ⇒ println("C: " + o), println, () ⇒ println("C Completed"))
    

    你只需要确保源 observable 是热的。最简单的方法是share它。

    【讨论】:

    • 如果你想让初始 observable 变冷怎么办?
    • @double_squeeze 只需使用publish 而不是share 并在所有订阅者都订阅时调用connect
    • 使用share 使其变热是没有意义的。提供的代码实际上订阅了 3 次 - 对于每个订阅者,与没有 share 的情况相同。 Krzysztof Skyrzynecki 的评论中描述了正确的做法:使用publish 而不是share,并在所有订阅者都订阅时调用connect
    • 我同意publish+connect 方法更简洁,我会在可行的情况下使用它(但有时,您只是事先不知道订阅者,也不在乎你错过了一些项目)。但是,我认为您关于三个订阅原始冷o 的说法是不正确的。如果您能证明其他情况,请share() :) 向我们提供您的证明。
    【解决方案2】:

    您不必从groupBy 折叠Observables。您可以改为订阅它们。

    类似这样的:

    String[] inputs= {"a", "b", "c", "a", "b", "b", "b", "a"};
    
    Action1<String> a = s -> System.out.print("-a-");
    
    Action1<String> b = s -> System.out.print("-b-");
    
    Action1<String> c = s -> System.out.print("-c-");
    
    Observable
        .from(inputs)
        .groupBy(s -> s)
        .subscribe((g) -> {
            if ("a".equals(g.getKey())) {
                g.subscribe(a);
            }
    
            if ("b".equals(g.getKey())) {
                g.subscribe(b);
            }
    
            if ("c".equals(g.getKey())) {
                g.subscribe(c);
            }
        });
    

    If 语句看起来有点难看,但至少您可以分别处理每个流。也许有办法避免它们。

    【讨论】:

    • 是的,我想尽可能避免那些 ifs。但是,如果它成功了,那么它看起来会更干净一些,因为它都在一个地方,而不是在原始流上进行过滤。谢谢!
    • 酷!如果我知道如何摆脱 if 语句,我会更新我的答案。
    • 您可以使用Dictionary&lt;groupKey,action&gt;,而您的组subscribe 方法必须从字典中解析操作并调用它。
    • Brandon Bill,我不得不重新阅读您的问题才能确定您最初使用的是filter。出于好奇,是什么让您更喜欢groupBy 解决方案而不是filter 解决方案?对我来说filter 似乎更简单,更容易理解,可能性能更好。你觉得groupBy 有什么优势?
    • 我在调试时注意到一些有趣的事情(注意这是使用 Rx JS),一旦“组”与处理程序相关联,组中元素的另一个实例将不会通过大订阅功能如上所示。相反,它将直接转发到该组订阅的 Action。然而,使用过滤器,可以对源 observable 的每个单个值进行比较。因此,对于 N 次拆分,将对每个元素进行可能的 N 次比较。对于组,只有在找到新组时才会进行比较。
    【解决方案3】:

    我一直在考虑这个问题,Tomas 解决方案还可以,但问题是它将流转换为热可观察对象。

    您可以将sharedefer 结合使用,以便与其他流一起获得冷可观察。

    例如(Java):

    var originalObservable = ...; // some source
    var coldObservable = Observable.defer(() -> {
        var shared - originalObservable.share();
        var aSource = shared.filter(x -> x.equals("a"));
        var bSource = shared.filter(x -> x.equals("b"));
        var cSource = shared.filter(x -> x.equals("c"));
        // some logic for sources
        return shared;
    });
    
    

    【讨论】:

      【解决方案4】:

      在 RxJava 中有一个特殊版本的 publish operator 接受一个函数。

      ObservableTransformer {
        it.publish { shared ->
          Observable.merge(
              shared.ofType(x).compose(transformerherex),
              shared.ofType(y).compose(transformerherey)
          )
        }
      }
      

      这会按类型拆分事件流。然后,您可以通过与不同的转换器组合来分别处理它们。他们都共享一个订阅。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2022-10-07
        相关资源
        最近更新 更多