Logo Questions Linux Laravel Mysql Ubuntu Git Menu
 

RXJava - combineLatest without losing any result

Tags:

java

rx-java

I want to combine two observables, one emits n items and the other one only 1.

combineLatest will wait until both observables have at least emitted one item and then combines the latest emitted items until both observable have finished. Consider following chronologically order:

  • Observable A -> emits result A1
  • Observable A -> emits result A2
  • Observable B -> emits result B1

combineLatest will only combine result 2 of observable 1 with result 1 of observable 2 (can be tester here easily: http://rxmarbles.com/#combineLatest).

What I need

I need to combine ALL items of two observables, no matter which one is faster. How can I do that?

Result should be (always, independent of which observable starts emitting items first!):

  • A1 combined with B1
  • A2 combined with B1
like image 876
prom85 Avatar asked Aug 15 '26 02:08

prom85


1 Answers

Old question, but I ran into this same issue. Here's my stab at it. First, the non-working version:

    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 );

Sure enough, the first couple emissions from emitsMany get dropped:

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

I think here's the fix.. First we need to wrap emitsOne into something that will continue to return the value previously observed without delay. I don't know of an operator that does this, but BehaviorSubject does the job.

Next we can use concatMap with a nested 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 );

We now get all the combinations:

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
like image 77
TrogDor Avatar answered Aug 17 '26 18:08

TrogDor



Donate For Us

If you love us? You can donate to us via Paypal or buy me a coffee so we can maintain and grow! Thank you!