RxJava中如何在计时器结束后重新触发Observable?
Solution: Chain Refresh with Timer in a Single Observable Stream
Instead of using doOnNext to spawn a separate interval stream, you can restructure your code to create a single, self-sustaining Observable chain using flatMap and timer. This keeps all operations connected and avoids separate, unlinked processes.
Step-by-Step Implementation
First, modify your refresh function to return an Observable (instead of being a void function) so we can chain operators seamlessly:
import io.reactivex.rxjava3.core.Observable import io.reactivex.rxjava3.schedulers.Schedulers import androidx.annotation.NonNull // Assume this is your response data class with the TTL field data class Response(val ttl: Long) fun refresh(): Observable<Response> { return service.makeNetworkCall() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .flatMap { response -> // Wait for the TTL to expire, then trigger the next refresh cycle Observable.timer(response.ttl, TimeUnit.SECONDS) .flatMap { refresh() } // Recursively call refresh to repeat the cycle .startWithItem(response) // Emit the current response immediately before waiting } }
How This Works
- Returning an Observable: By making
refresh()return anObservable<Response>, we can chain all operations into one continuous sequence. - flatMap for Chaining: After receiving the network response,
flatMapswitches to atimerthat waits the specified TTL. Once the timer completes, we recursively callrefresh()to start the entire cycle again. - startWithItem: This ensures you get the current response right away (just like your original code) before waiting for the TTL to expire—no delay in using the latest data.
- Single Stream Control: Unlike your previous
interval+doOnNextapproach, this creates a single Observable chain. Disposing the subscription will cleanly stop all future refreshes, with no orphaned streams running in the background.
Using the Observable
To start the refresh cycle and manage its lifecycle:
private var refreshDisposable: Disposable? = null fun startRefreshing() { refreshDisposable = refresh() .subscribe( { response -> // Handle the latest response (update UI, cache data, etc.) println("Received response with TTL: ${response.ttl} seconds") }, { error -> // Handle network errors (optional: add retry logic or stop the cycle) error.printStackTrace() } ) } fun stopRefreshing() { refreshDisposable?.dispose() refreshDisposable = null }
Optional: Error Handling
If you want to retry on network failures, add operators like retry to the chain. For example, to retry up to 3 times before stopping:
fun refresh(): Observable<Response> { return service.makeNetworkCall() .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .retry(3) // Retry 3 times on error before propagating the failure .flatMap { response -> Observable.timer(response.ttl, TimeUnit.SECONDS) .flatMap { refresh() } .startWithItem(response) } }
Key Benefits Over Your Original Approach
- No Unlinked Streams: Everything is part of one Observable chain, giving you full control over lifecycle and cleanup.
- Dynamic TTL Adaptation: Each refresh uses the latest TTL from the response—if the TTL changes between calls, the cycle automatically adjusts (your interval approach would stick to the initial TTL value).
- Cleaner, Maintainable Code: The logic is self-contained in the chain, making it easier to follow and modify later.
内容的提问来源于stack exchange,提问作者2bornot2b
相关产品推荐
相关产品推荐

