老问题,但我遇到了同样的问题。这是我的尝试。一、非工作版:
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