【问题标题】:Rxjava : takeUntil skip complete eventRxjava : takeUntil 跳过完成事件
【发布时间】:2020-11-05 07:36:53
【问题描述】:

TakeUntil
在第二个 Observable 发出一个项目或终止后丢弃任何由 Observable 发出的项目

有没有办法在使用takeUntil 运算符时跳过complete 事件?

Observable<Long> publish = source.publish(
        multicast -> multicast.flatMapMaybe(
                o -> Observable.interval(0, 500, TimeUnit.MILLISECONDS, scheduler)
                        .takeUntil(multicast)
                        .reduce(Long::sum)
        )
);

此图说明了上面代码的生成结果:

//events: --------x-----1----2---1---x-----3--0--------x-1---1----|  
//result: ---------------------------4-----------------3----------2  

我想要处理 4 和 3 事件,这与上一个事件 2 不同,后者是由 complete 的 source 事件生成的。

【问题讨论】:

    标签: rx-java rx-java2


    【解决方案1】:

    本期为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);
    }
    

    【讨论】:

    • 再次感谢,但我不明白 :( ... 为什么有两个 takeUntil 以及为什么我应该使用 switchIfEmpty (这是我的代码的简化示例,我使用不可变对象而不是比long) ? 我知道这是一个复杂的用例:)
    • 我可能不明白,你想要做什么。您是想计算“最后一个”值 2 并在以后以不同方式处理它,还是想完全丢弃“最后一个”值的 reduce 并发出不同的东西?
    • 我想丢弃 le 最后一个值的 reduce ... 是的,我认为这不可能我需要其他东西来过滤
    • 如果你不想处理最后发出的“reduce”的结果,那么我的解决方案就足够了。我将更新答案以详细解释实际发生的情况。
    猜你喜欢
    • 2015-09-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-06-17
    相关资源
    最近更新 更多