本期为64676462的后续
我的结果并不漂亮,但我认为它确实有效。您的用例并不容易解决,因为您不知道发出的值是“最后一个”,因为它是一个流。您不能发出“最后一个”值并知道它是最后一个,因为发出和 onComplete 是两个原子操作。您需要延迟发出“最后一个”值并等待 onComplete 事件,以确保“最后一个”值确实是最后一个值。
你也可以作弊,就像我在这里做的那样:
注意:当源 observable 和内部订阅 observable 完成时,flatMap 完成。因此,我们必须确保当源完成时,内部流也完成。这是通过使用另一个#takeUntil 来实现的,它通过#materialize 从源侦听onComplete-event。 reduce 之后的 #takeUntil 确保不会向下游发出 reduce 值。最后 #switchIfEmpty 会将 onComplete 事件从“last”值转换为另一个值,因为 observable 没有发出任何值。
注意:假设是,所有值都从一个线程发出同步。
@Test
public void takeWhileReduce() {
TestScheduler scheduler = new TestScheduler();
PublishSubject<Integer> source = PublishSubject.create();
Observable<Long> publish = source.publish(
multicast -> {
return multicast.flatMap(
o -> {
return Observable.interval(0, 500, TimeUnit.MILLISECONDS, scheduler) //
.takeUntil(multicast)
.reduce(Long::sum)
.toObservable()
// make sure, that the inner stream completes, when the outer stream completes.
// takeUntil must be after reduce, because takeUntil will close the stream and therefore reduce
// will push its value to the subscriber.
.takeUntil(multicast.materialize().filter(Notification::isOnComplete))
// when the upstream is closed, switch over to a fallback observable.
// if you want special handling for the "LAST" value, just provide another fallback observable.
.switchIfEmpty(Observable.just(Long.MAX_VALUE));
},
1);
});
TestObserver<Long> test = publish.test();
source.onNext(42);
scheduler.advanceTimeBy(1500, TimeUnit.MILLISECONDS);
// action - push next value - flatMapped value will complete and push value
source.onNext(42);
// assert - values emitted: 0,1,2,3
test.assertValuesOnly(6L);
// next value is flatMapped
scheduler.advanceTimeBy(1000, TimeUnit.MILLISECONDS);
// action - push next value - flatMapped value will complete and push value
source.onNext(42);
// assert - values emitted: 0,1,2
test.assertValuesOnly(6L, 3L);
scheduler.advanceTimeBy(500, TimeUnit.MILLISECONDS);
// action - push next value - flatMapped value will complete and push value
source.onNext(42);
// assert - values emitted: 0,1
test.assertValuesOnly(6L, 3L, 1L);
scheduler.advanceTimeBy(500, TimeUnit.MILLISECONDS);
// action - outer-stream completes
source.onComplete();
test.assertComplete().assertValues(6L, 3L, 1L, Long.MAX_VALUE);
}