RxJava轮询接口时请求超时被取消,如何避免连续无响应?
The problem with your current setup is that Observable.interval emits a new tick every second, regardless of whether your previous Retrofit request has completed. switchMap cancels the ongoing inner observable (your Retrofit call) as soon as a new tick arrives—this is exactly what it's designed to do, but it's not what you want here when requests can take longer than your interval.
To fix this, we need to adjust the polling logic so that we never cancel an ongoing request, and instead start the next request either:
- 1 second after the previous request started (if it completed quickly), or
- Immediately after the previous request finishes (if it took longer than 1 second)
Here's a robust implementation that achieves this:
import io.reactivex.Observable; import io.reactivex.schedulers.Schedulers; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; public Observable<YourResponseType> createPollingObservable() { return Observable.defer(() -> { // Track when each poll starts to maintain our 1-second interval AtomicLong lastPollStartTime = new AtomicLong(System.currentTimeMillis()); return Observable.just(0) // Repeat after calculating the appropriate delay .repeatWhen(completed -> completed.flatMap(ignored -> { long currentTime = System.currentTimeMillis(); long timeSinceLastStart = currentTime - lastPollStartTime.get(); // Calculate how long to wait before starting the next poll long delayMs = Math.max(0, 1000 - timeSinceLastStart); return Observable.timer(delayMs, TimeUnit.MILLISECONDS, Schedulers.io()); })) .flatMap(ignored -> { // Update the start time before launching the request lastPollStartTime.set(System.currentTimeMillis()); return pollWithRetrofit(); // Your existing Retrofit call observable }); }); }
How this works:
Observable.defer: Ensures each subscription gets its ownlastPollStartTimeinstance, avoiding shared state issues.repeatWhen: Instead of using a fixed interval, we wait until the current poll completes, then calculate how much time we need to wait to get back to our 1-second cadence. If the previous request took longer than 1 second, we start the next one immediately (delay of 0).- No cancellation: Since we only trigger the next poll after the previous one completes, there's never a case where an ongoing request is canceled by a new tick.
Alternative Simplified Approach
If you don't need to strictly maintain a 1-second interval between the start of requests (and just want to wait 1 second after each response before starting the next), you can use this simpler version:
public Observable<YourResponseType> createSimplePollingObservable() { return Observable.defer(this::pollWithRetrofit) .repeatWhen(completed -> completed.delay(1, TimeUnit.SECONDS, Schedulers.io())) .subscribeOn(Schedulers.io()); }
This will start a new request exactly 1 second after the previous one finishes. It's less strict on the interval but still ensures no requests are canceled.
Why Your Original Code Fails
Observable.interval emits every second unconditionally. When switchMap receives a new tick, it cancels the previous Retrofit observable (which is still running if the request took longer than 1 second). This can lead to a situation where you're constantly canceling requests before they complete, resulting in no responses being emitted for long periods.
By switching to a completion-driven repeat pattern, we eliminate this cancellation behavior entirely.
内容的提问来源于stack exchange,提问作者TpoM6oH

