【问题标题】:RxJava `Completable.andThen` is not executing serially?RxJava `Completable.andThen` 没有连续执行?
【发布时间】:2018-02-05 14:32:17
【问题描述】:

我有一个用例,我在 Completable 中初始化一些全局变量,然后在链的下一步(使用 andThen 运算符)中使用这些变量。

以下示例详细解释了我的用例

假设你有一堂课User

        class User {
            String name;
        }

我有一个像这样的 Observable,

        private User mUser; // this is a global variable

        public Observable<String> stringObservable() {
            return Completable.fromAction(() -> {
                mUser = new User();
                mUser.name = "Name";
            }).andThen(Observable.just(mUser.name));
        }           

首先,我在 Completable.fromAction 中进行一些初始化,我希望 andThen 运算符仅在完成 ​​Completable.fromAction 后启动。

这意味着我希望mUser 在andThen 运算符启动时被初始化。

以下是我对这个 observable 的订阅

             stringObservable()
            .subscribe(s -> Log.d(TAG, "success: " + s),
                    throwable -> Log.e(TAG, "error: " + throwable.getMessage()));

但是当我运行这段代码时,我得到一个错误

          Attempt to read from field 'java.lang.String User.name' on a null object reference

这意味着mUser 为空,andThen 在执行Completable.fromAction 中的代码之前启动。这里发生了什么事?

根据andThen的文档

返回一个 Observable,它将订阅这个 Completable,一旦完成,就会订阅 {@code next} ObservableSource。来自此 Completable 的错误事件将传播到下游订阅者,并导致跳过 Observable 的订阅。

【问题讨论】:

    标签: java rx-java rx-java2


    【解决方案1】:

    问题不在于andThen,而在于andThen 中的Observable.just(mUser.name) 语句。 just 运算符将尝试立即创建 observable,尽管它只会在 Completable.fromAction 之后发出。

    这里的问题是,在尝试使用 just 创建 Observable 时,mUser 为空。

    解决方案:您需要推迟创建 String Observable 直到订阅发生,直到 andThen 的上游开始发射。

    而不是andThen(Observable.just(mUser.name));

    使用

     andThen(Observable.defer(() -> Observable.just(mUser.name)));
    

    或者

     andThen(Observable.fromCallable(() -> mUser.name));
    

    【讨论】:

    • Observable.just(T) docs 声明:“请注意,该项目是按原样获取并重新发送,而不是通过任何方式计算。使用fromCallable(Callable) 到按需生成单个项目(当观察者订阅它时)。”
    • 是的。答案中也提到了
    • 在答案中直接引用文档是一个不错的选择,但是
    【解决方案2】:

    我不认为@Sarath Kn 的回答是 100% 正确的。是的,just 会在调用后立即创建 observable,但andThen 仍在意外时间调用just。

    我们可以比较 andThen 和 flatMap 以获得更好的理解。这是一个完全可运行的测试:

    package com.example;
    
    import org.junit.Test;
    
    import io.reactivex.Completable;
    import io.reactivex.Observable;
    import io.reactivex.observers.TestObserver;
    import io.reactivex.schedulers.Schedulers;
    
    public class ExampleTest {
    
        @Test
        public void createsIntermediateObservable_AfterSubscribing() {
            Observable<String> coldObservable = getObservableSource()
                    .flatMap(integer -> getIntermediateObservable())
                    .subscribeOn(Schedulers.trampoline())
                    .observeOn(Schedulers.trampoline());
            System.out.println("Cold obs created... subscribing");
            TestObserver<String> testObserver = coldObservable.test();
            testObserver.awaitTerminalEvent();
    
            /*
            Resulting logs:
    
            Creating observable source
            Cold obs created... subscribing
            Emitting 1,2,3
            Creating intermediate observable
            Creating intermediate observable
            Creating intermediate observable
            Emitting complete notification
    
            IMPORTANT: see that intermediate observables are created AFTER subscribing
             */
        }
    
        @Test
        public void createsIntermediateObservable_BeforeSubscribing() {
            Observable<String> coldObservable = getCompletableSource()
                    .andThen(getIntermediateObservable())
                    .subscribeOn(Schedulers.trampoline())
                    .observeOn(Schedulers.trampoline());
            System.out.println("Cold obs created... subscribing");
            TestObserver<String> testObserver = coldObservable.test();
            testObserver.awaitTerminalEvent();
    
            /*
            Resulting logs:
    
            Creating completable source
            Creating intermediate observable
            Cold obs created... subscribing
            Emitting complete notification
    
            IMPORTANT: see that intermediate observable is created BEFORE subscribing =(
             */
        }
    
        private Observable<Integer> getObservableSource() {
            System.out.println("Creating observable source");
            return Observable.create(emitter -> {
                System.out.println("Emitting 1,2,3");
                emitter.onNext(1);
                emitter.onNext(2);
                emitter.onNext(3);
                System.out.println("Emitting complete notification");
                emitter.onComplete();
            });
        }
    
        private Observable<String> getIntermediateObservable() {
            System.out.println("Creating intermediate observable");
            return Observable.just("A");
        }
    
        private Completable getCompletableSource() {
            System.out.println("Creating completable source");
            return Completable.create(emitter -> {
                System.out.println("Emitting complete notification");
                emitter.onComplete();
            });
        }
    }
    

    您可以看到,当我们使用flatmap 时,just 在 订阅之后被调用,这是有道理的。如果中间可观察对象依赖于发送到flatmap 的项目,那么系统当然无法在订阅之前创建中间可观察对象。它还没有任何值。你可以想象如果flatmap在订阅之前调用just这不会起作用:

    .flatMap(integer -> getIntermediateObservable(integer))
    

    奇怪的是andThen 能够在订阅之前创建它的内部可观察对象(即调用just)。它可以做到这一点是有道理的。 andThen 唯一会收到一个完整的通知,因此没有理由不尽早创建中间可观察对象。唯一的问题是它不是预期的行为。

    @Sarath Kn 的解决方案是正确的,但原因是错误的。如果我们使用defer,我们可以看到一切正常:

    @Test
    public void usingDefer_CreatesIntermediateObservable_AfterSubscribing() {
        Observable<String> coldObservable = getCompletableSource()
                .andThen(Observable.defer(this::getIntermediateObservable))
                .subscribeOn(Schedulers.trampoline())
                .observeOn(Schedulers.trampoline());
        System.out.println("Cold obs created... subscribing");
        TestObserver<String> testObserver = coldObservable.test();
        testObserver.awaitTerminalEvent();
    
        /*
        Resulting logs:
    
        Creating completable source
        Cold obs created... subscribing
        Emitting complete notification
        Creating intermediate observable
    
        IMPORTANT: see that intermediate observable is created AFTER subscribing =) YEAY!!
         */
    }
    

    【讨论】:

    • 我认为这个答案是正确的,但很长。我们只需要知道andThen 接受一个可完成项并返回它(或接受一个并返回它)。如果你做func1(func2()),func2 将不会在func1 的作用域内被调用,它会被称为调用者作用域。莎拉将此与.just 联系起来是错误的。
    猜你喜欢
    • 2021-01-26
    • 2018-10-19
    • 1970-01-01
    • 2021-11-16
    • 1970-01-01
    • 1970-01-01
    • 2012-01-29
    • 2015-09-12
    • 1970-01-01
    相关资源
    最近更新 更多