【问题标题】:RXJava - combineLatest without losing any resultRXJava - combineLatest 不丢失任何结果
【发布时间】:2016-09-07 17:53:20
【问题描述】:

我想组合两个 observable,一个发出 n 个项目,另一个只发出 1 个。

combineLatest 将等到两个 observable 至少发出一个项目,然后合并最新发出的项目,直到两个 observable 都完成。请考虑按时间顺序:

  • Observable A -> 发出结果 A1
  • Observable A -> 发出结果 A2
  • Observable B -> 发出结果 B1

combineLatest 只会将 observable 1 的结果 2 与 observable 2 的结果 1 结合起来(可以在这里轻松地成为测试人员:http://rxmarbles.com/#combineLatest)。

我需要什么

我需要结合两个可观察对象的所有项目,无论哪一个更快。我该怎么做?

结果应该是(总是,独立于哪个 observable 先开始发射项目!):

  • A1 与 B1 结合
  • A2 与 B1 结合

【问题讨论】:

  • 你的意思是你想要一些序列a的所有排列,序列b?或者你想压缩它们? reactivex.io/documentation/operators/zip.html
  • 所有排列...一个 observable 只发出 1 项,其他 n 项,我想得到 n 个组合(所有 n 项,每个项与 1 项组合)
  • 我调整了我的示例以使其更清晰...
  • 可以尝试在 Observable 2 中链接 startWith 运算符,这样您就不会丢失任何值。然后你会得到,A1 - Empty, A2- B1.. 然后所有后续排放将与其他可观察到的最新排放合并。

标签: java rx-java


【解决方案1】:

老问题,但我遇到了同样的问题。这是我的尝试。一、非工作版:

    Observable<Integer> emitsMany = Observable.range( 1, 10 )
            .concatMap( i -> Observable.just( i ).delay( 1, TimeUnit.SECONDS ))
            .doOnNext( i -> System.out.println( "produced " + i ));

    Observable<Boolean> emitsOne = Observable.just( true )
            .delay( 3, TimeUnit.SECONDS )
            .doOnNext( b -> System.out.println( "produced " + b ));

    Observable.combineLatest(
            emitsMany, emitsOne,
            ( i, b ) -> "consumed " + i + " " + b )
    .blockingSubscribe( System.out::println );

果然,emitsMany 的前几个排放被删除了:

produced 1
produced 2
produced 3
produced true
consumed 3 true
produced 4
consumed 4 true
. . .

我认为这是解决方法。首先,我们需要将emitsOne 包装成将继续立即返回先前观察到的值的东西。我不知道有哪个运营商可以做到这一点,但 BehaviorSubject 可以做到这一点。

接下来我们可以将concatMap与嵌套的take(1) Observable一起使用:

    Observable<Integer> emitsMany = Observable.range( 1, 10 )
            .concatMap( i -> Observable.just( i ).delay( 1, TimeUnit.SECONDS ))
            .doOnNext( i -> System.out.println( "produced " + i ));

    Observable<Boolean> emitsOne = Observable.just( true )
            .delay( 3, TimeUnit.SECONDS )
            .doOnNext( b -> System.out.println( "produced " + b ));

    BehaviorSubject<Boolean> emitsOneSubject = BehaviorSubject.create();
    emitsOne.subscribe( emitsOneSubject::onNext );

    emitsMany.concatMap( i -> emitsOneSubject
            .take( 1 )
            .map( b -> "consumed " + i + " " + b ))
    .blockingSubscribe( System.out::println );

我们现在得到了所有的组合:

produced 1
produced 2
produced true
consumed 1 true
consumed 2 true
produced 3
consumed 3 true
produced 4
consumed 4 true
produced 5
consumed 5 true
produced 6
consumed 6 true
produced 7
consumed 7 true
produced 8
consumed 8 true
produced 9
consumed 9 true
produced 10
consumed 10 true

【讨论】:

    【解决方案2】:

    请注意,这是未经测试的:

    从可观察序列a 创建一个ReplaySubject。对于在序列b 上发出的每个值,将该值与重播主题相结合,以创建一个新的 Pair&lt;A, B&gt; 可观察对象。将这些 observable 平面映射在一起并将其作为结果返回。

    public static <A, B> Observable<Pair<A, B>> permutation(
        Observable<A> observableA, 
        Observable<B> observableB, 
    ) {
        ReplaySubject<A> subjectA = ReplaySubject.create();
        observableA.subscribe(subjectA::onNext);
        return observableB.flatMap(b -> subjectA.map(a -> Pair.of(a, b)));
    }
    

    【讨论】:

    • 我会试试这个...我实际上想要一个支持顺序执行任意数量的项目的函数...我为此编写了一个有效的函数,但它并不漂亮而不是真正的 rx java...
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多