【问题标题】:RxJava: merge() changes the order of the emitted items?RxJava:merge() 改变了发射项目的顺序?
【发布时间】:2018-12-28 14:48:20
【问题描述】:

看看这两个小测试:

@Test
public void test1() {
    Observable.range(1, 10)
        .groupBy(v -> v % 2 == 0)
        .flatMap(group -> {
            if (group.getKey()) {
                return group;
            }
            return group;
        })
        .subscribe(System.out::println);
}

@Test
public void test2() {
    Observable.range(1, 10)
        .groupBy(v -> v % 2 == 0)
        .toMap(g -> g.getKey())
        .flatMapObservable(m -> Observable.merge(
            m.get(true),
            m.get(false)))
        .subscribe(System.out::println);
}

我希望两者都以相同的顺序返回一个数字列表,所以:

1 2 3 4 5 6 7 8 9 10

但第二个例子返回

2 4 6 8 10 1 3 5 7 9

改为。

似乎在第二个示例中merge 正在执行concat,实际上如果我将其更改为concat,结果是相同的。

我错过了什么?

谢谢。

【问题讨论】:

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


    【解决方案1】:

    基本上flatMapmerge不保证发射项目的顺序。

    来自flatMap文档:

    请注意,FlatMap 会合并这些 Observable 的发射,以便它们可以交错。

    来自merge文档:

    Merge 可以交错由合并后的 Observable 发出的项(类似的运算符 Concat 不会交错项,而是在开始从下一个源 Observable 发出项之前依次发出每个源 Observable 的所有项)。

    引用此SO Answer

    在您的情况下,对于单元素静态流,它没有任何真正的区别(但理论上,合并可以以随机顺序输出单词并且仍然根据规范有效)

    如果您需要保证订单,请改用concat*

    第一个例子

    它是这样工作的:

    • 当发出1 时,groupBy 运算符将创建一个GroupedObservable,其键为false
      • flatMap 将从这个 observable 输出项目 - 目前只有 1
    • 当发出2 时,groupBy 运算符将创建一个GroupedObservable,其键为true
      • flatMap 现在也将输出第二个 GroupedObservable 的项目 - 目前是 2
    • 当发出3 时,groupBy 运算符会将其添加到现有的GroupedObservable 中,键为falseflatMap 将立即输出此项目
    • 当发出4 时,groupBy 运算符会将其添加到现有的 GroupedObservable 中,键为 trueflatMap 将立即输出此项目

    它可能会帮助您添加更多日志记录:

        Observable.range(1, 10)
                .groupBy(v -> v % 2 == 0)
                .doOnNext(group -> System.out.println("key: " + group.getKey()))
                .flatMap(group -> {
                    if (group.getKey()) {
                        return group;
                    }
                    return group;
                })
                .subscribe(System.out::println);
    

    那么输出是:

    key: false
    1
    key: true
    2
    3
    ...
    

    第二个例子

    这是完全不同的,因为toMap 将阻塞直到上游完成:

    • 当发出1 时,groupBy 运算符将创建一个GroupedObservable,其键为false
      • toMap 将把这个GroupedObservable 添加到内部映射并使用键false(与GroupedObservable 具有相同的键)
    • 当发出2 时,groupBy 运算符将创建一个GroupedObservable,其键为true
      • toMap 将把这个GroupedObservable 添加到内部映射并使用键true(与GroupedObservable 具有相同的键) - 所以现在映射有2 个GroupedObservables
    • 以下数字被添加到相应的GroupedObservables 中,当源代码完成时,toMap 运算符完成并将映射传递给下一个运算符
    • flatMapObservable 中,您使用映射创建一个新的可观察对象,您首先添加偶数元素(键 = true),然后添加奇数元素(键 = false

    您还可以在此处添加更多日志记录:

        Observable.range(1, 10)
                .groupBy(v -> v % 2 == 0)
                .doOnNext(group -> System.out.println("key: " + group.getKey()))
                .toMap(g -> g.getKey())
                .doOnSuccess(map -> System.out.println("map: " + map.size()))
                .flatMapObservable(m -> Observable.merge(
                        m.get(true),
                        m.get(false)
                ))
                .subscribe(System.out::println);
    

    那么输出是:

    key: false
    key: true
    map: 2
    2
    4
    6
    8
    10
    1
    3
    5
    7
    9
    

    【讨论】:

    • 感谢您的回答,但我仍然不明白为什么 flatMapObservable 内的 merge 不保留订单。我现在知道toMap 正在阻塞,因为它需要知道地图的键才能在下一阶段发出它。 In your case, with single-element, static streams, it is not making any real difference 但新地图中的 2 Observable 不是单元素,是吗?我也不确定static 在这种情况下是什么意思。 merge 的定义表明元素可能交错,这是我在这种情况下所期望的,而是表现得像 concat
    • 对,地图中的 2 个 Observable 各有 5 个元素 - 但它们是 static(即不依赖于某些时间) - 换句话说:流已经完全定义,所有元素都是已知。 toMap 源流已经完成,map 中的 Observable 订阅后可以立即输出所有元素。
    • ad 2nd comment:好吧,在这种情况下,merge 的行为与concat 的行为相同只是巧合。所以这个merge 实现显然订阅了第一个可观察对象(您作为第一个参数传递的内容)并处理其所有数据,然后订阅第二个并处理第二个可观察对象的所有数据。但如顶部所述:merge 无法保证这一点。
    • 好吧,巧合可能是错误的词。但我想表达的是,我们看到的merge 的行为是依赖于实现的。换句话说:其他 Rx 实现可能有不同的行为——甚至更新版本的 RxJava 也可能有不同的行为。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-12-02
    • 2016-12-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-20
    相关资源
    最近更新 更多