RxJava订阅成功后仍抛出TimeoutException?问题咨询与解决
My Code & Problem
I wrote the following RxJava code:
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()
The issue is: I get the "Success" log within a few seconds, but my loader keeps waiting for a full minute, after which the error callback of the subscription is executed. Is this expected behavior? How can I stop the timeout mechanism after a successful subscription?
Additional Note: The Observable returned by getObservable() is created via PublishSubject.create().
Answer
This is not expected behavior—the problem stems from how PublishSubject works and how you've structured your Rx chain. Let's break this down:
Why This Happens
PublishSubject is a hot Observable that emits events to all current subscribers, but it doesn't automatically terminate on its own. The timeout operator's logic is: if no event (including onComplete) is emitted within the specified window, trigger a timeout error.
In your case, the PublishSubject emits a success event but never calls onComplete(). So the timeout operator keeps waiting for the next event, even after the success callback runs. After 1 minute hits, it triggers the timeout error, which fires your error callback and runs doFinally again (resetting your loader state once more).
Fixes You Can Implement
1. Call onComplete() in the PublishSubject
If this Observable is only meant to emit a single event, add onComplete() right after emitting your success data. This tells the Observable (and timeout) to terminate cleanly:
// Inside the code that creates your PublishSubject val subject = PublishSubject.create<YourDataModel>() subject.onNext(yourSuccessData) subject.onComplete() // Add this line to end the stream
2. Add take(1) to your chain
If you can't modify the getObservable() implementation, add take(1) to force the stream to terminate after the first emission. This will stop timeout from waiting for more events:
repo.getObservable() .take(1) // Terminate after first emission .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()
3. Bonus: Understand Hot vs Cold Observables
Hot Observables like PublishSubject don't have a built-in termination point unless you explicitly call onComplete() or onError(). If you don't need a hot stream, consider using a cold Observable (like Observable.create() that emits once and completes) instead—this would avoid the issue entirely.
内容的提问来源于stack exchange,提问作者SMGhost

