【问题标题】:Rx: How to get the last element even if onError was called?Rx:即使调用了 onError,如何获取最后一个元素?
【发布时间】:2016-06-09 19:18:09
【问题描述】:

我正在使用 RxJava,我需要做两件事:

  • 获取从Observable 发出的最后一个元素
  • 确定是否调用了onError,与onCompleted

我已经研究过使用lastlastOrDefault(这实际上是我需要的行为),但我无法解决onError 隐藏最后一个元素的问题。我可以使用 Observable 两次,一次获取 last 值,一次获取完成状态,但到目前为止,我只能通过创建自己的 Observer 来完成此操作:

public class CacheLastObserver<T> implements Observer<T> {

    private final AtomicReference<T> lastMessageReceived = new AtomicReference<>();
    private final AtomicReference<Throwable> error = new AtomicReference<>();

    @Override
    public void onCompleted() {
        // Do nothing
    }

    @Override
    public void onError(Throwable e) {
        error.set(e);
    }

    @Override
    public void onNext(T message) {
        lastMessageReceived.set(message);
    }

    public Optional<T> getLastMessageReceived() {
        return Optional.ofNullable(lastMessageReceived.get());
    }

    public Optional<Throwable> getError() {
        return Optional.ofNullable(error.get());
    }
}

我自己创建Observer 没有问题,但感觉 Rx 应该能够更好地满足“在完成之前获取最后一个元素发出”的用例。关于如何实现这一点的任何想法?

【问题讨论】:

    标签: java rx-java reactivex


    【解决方案1】:

    试试这个:

    source.materialize().buffer(2).last()
    

    在错误情况下,最后一个发射将是两个项目的列表,即最后发射的值包装为Notification 和错误通知。如果没有错误,第二项将是完成通知。

    还要注意,如果没有发出任何值,那么结果将是一个列表,其中包含一项作为终端通知。

    【讨论】:

    • 有没有办法获取最后一个发射的项目,而不是最后收到的项目,这是一个很好的方法,虽然在 onNext 调用中获取最后一个项目,但不是从初始 Observable 的发射. (即发射,1,2,3,在 4 之前崩溃,那么我应该有崩溃和第 4 项,目前您将使用此方法获得 3)。
    • 源发出 1,2,3, error, 4 违反 Observable 合约,因此 RxJava 不支持。您必须抑制错误的抛出并将其作为 onNext 发射来实现。
    • 这不适用于具有奇数项的可观察对象(完整/错误通知在单例列表中)。即使source.materialize().buffer(2, 1).last() 也无法解决,因为最后一个缓冲区的长度为 1。
    【解决方案2】:

    我解决了:

    source.materialize().withPrevious().last()
    

    withPrevious 在哪里 (Kotlin):

    fun <T> Observable<T>.withPrevious(): Observable<Pair<T?, T>> =
        this.scan(Pair<T?, T?>(null, null)) { previous, current -> Pair(previous.second, current) }
            .skip(1)
            .map { it as Pair<T?, T> }
    

    【讨论】:

      【解决方案3】:

      有没有尝试过onErrorResumeNext 这里可以看到其余的或者错误处理操作符https://github.com/ReactiveX/RxJava/wiki/Error-Handling-Operators

      【讨论】:

        【解决方案4】:

        我使用这种方法来解决您的问题。

        public class ExampleUnitTest {
            @Test
            public void testSample() throws Exception {
                Observable.just(1, 2, 3, 4, 5)
                        .map(number -> {
                            if (number == 4)
                                throw new NullPointerException();
                            else
                                return number;
                        })
                        .onErrorResumeNext(t -> Observable.empty())
                        .lastOrDefault(15)
                        .subscribe(lastEmittedNumber -> System.out.println("onNext: " + lastEmittedNumber));
            }
        }
        

        它会发出onNext: 3

        希望对你有帮助。

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 2020-04-06
          • 1970-01-01
          • 1970-01-01
          • 2021-06-19
          • 1970-01-01
          • 1970-01-01
          • 2019-11-16
          • 2020-03-28
          相关资源
          最近更新 更多