You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RxJava订阅成功后仍抛出TimeoutException?问题咨询与解决

RxJava: Timeout triggers error even after successful emission?

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.14 09:15:32