【问题标题】:rxjava timer throws timeout exception after subscribe has successfully executed订阅成功执行后 rxjava 计时器抛出超时异常
【发布时间】:2019-12-19 08:20:46
【问题描述】:

我有以下代码:

repo.getObservable()
            .timeout(1, TimeUnit.MINUTES)
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .doOnSubscribe {
                _isInProgress.value = true
            }
            .doFinally {
                _isInProgress.value = false
            }
            .subscribe(
                    {
                        Timber.d("Success")
                    },
                    {

                        Timber.e(it)
                    })
            .trackDisposable()

它的问题是几秒钟后我成功地收到了 Success 消息,但我的预加载器仍然等待 1 分钟,然后我的订阅错误部分被执行。这是预期的行为吗?如果订阅的成功部分被执行,我该怎么做才能停止超时?

P。 S. getObservable() 返回的 Observable 是这样创建的:PublishSubject.create()

【问题讨论】:

    标签: android rx-java rx-java2


    【解决方案1】:

    如果您需要一个结果,请在timeout 之前或之后使用take(1)

    repo.getObservable()
            .take(1) // <---------------------------------------------
            .timeout(1, TimeUnit.MINUTES)
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .doOnSubscribe {
                _isInProgress.value = true
            }
            .doFinally {
                _isInProgress.value = false
            }
            .subscribe(
                    {
                        Timber.d("Success")
                    },
                    {
    
                        Timber.e(it)
                    })
            .trackDisposable()
    

    【讨论】:

      【解决方案2】:

      那一定是因为您没有在 PublishSubject 上调用 onComplete,而只调用了 onNext。 看这个例子:

      PublishSubject<Integer> source = PublishSubject.create();
      
      // It will get 1, 2, 3, 4 and onComplete
      source.subscribe(getFirstObserver()); 
      
      source.onNext(1);
      source.onNext(2);
      source.onNext(3);
      
      // It will get 4 and onComplete for second observer also.
      source.subscribe(getSecondObserver());
      
      source.onNext(4);
      source.onComplete();
      

      直到onComplete 被调用,观察者都在等待更多结果。 您可以在收到您等待的结果后unsubscribe/dispose,也可以在发送完所有结果后在 Observable 上调用onComplete

      【讨论】:

      • 我从不调用 onComplete,因为我的主题应该是可重用的。得到结果后如何退订?如何在订阅方法中访问一次性?
      • repo.getObservable()返回值保存在一个变量中,我建议您只在成功回调中调用一个方法,这样您就可以在该方法中调用您当前的代码(Timber.d("Success"))以及取消订阅先前保存的 var 的方法。请记住,如果该 observable 仍然处于活动状态并且 Activity 或 Fragment 被破坏或应用程序将出现内存泄漏,您还应该取消订阅该 observable。
      • getObservable() 返回 PublishSubject。如果我将其保存到变量中,我该如何使用它来取消订阅? PublishSubject 上没有取消订阅方法。
      • PublishSubjectSubject 类型,它是 Observable 类型。而ObservableunsubscribeOn(scheduler) 方法。而你设置的调度器是AndroidSchedulers.mainThread()
      • 是的,有 unsubscribeOn 方法,但它所做的只是指定应该在哪个线程上执行取消订阅操作。它没有给我一次性对象,我可以在“订阅成功”方法中使用它来取消订阅。
      猜你喜欢
      • 2016-12-11
      • 1970-01-01
      • 2020-05-01
      • 1970-01-01
      • 2021-03-31
      • 2021-03-25
      • 1970-01-01
      • 1970-01-01
      • 2021-12-18
      相关资源
      最近更新 更多