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

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

  1. Returning an Observable: By making refresh() return an Observable<Response>, we can chain all operations into one continuous sequence.
  2. flatMap for Chaining: After receiving the network response, flatMap switches to a timer that waits the specified TTL. Once the timer completes, we recursively call refresh() to start the entire cycle again.
  3. 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.
  4. Single Stream Control: Unlike your previous interval + doOnNext approach, 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 10:05:25