RxJava2数据轮询实现求助:RetryWhen与单订阅处理器相关问题
Hey there! Let's walk through how to implement data polling in RxJava 2 with your existing getMyTask() Single, and tackle your questions about retryWhen and single-subscriber processors.
Data polling in RxJava boils down to re-running your task after a delay—whether it succeeds or fails, depending on your requirements. We'll use retryWhen (for handling retries after failures) and repeatWhen (for triggering new runs after successful completion) to build this flow.
1. Using retryWhen for Retries on Failure
The retryWhen operator lets you control when to retry a failed Observable. It takes an Observable<Throwable> (emitted when your task fails) and returns an Observable: if this returned Observable emits an onNext event, the task retries; if it emits onComplete or onError, the flow stops.
Here's how to apply it to your getMyTask() to retry after 3 seconds when your custom IllegalArgumentException is thrown:
getMyTask() .retryWhen(throwables -> throwables.flatMap(throwable -> { // Only retry for our specific error type if (throwable instanceof IllegalArgumentException) { Log.d("ERSEN", "Retrying task after error..."); return Observable.timer(3, TimeUnit.SECONDS); } // Terminate the flow for other errors return Observable.error(throwable); })) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe( result -> Log.d("ERSEN", "Task succeeded: " + result), error -> Log.d("ERSEN", "Flow terminated with error: " + error.getMessage()) );
2. Continuous Polling (Success + Failure Handling)
If you want to keep polling even after the task succeeds (wait a delay, then run again), combine retryWhen with repeatWhen. The repeatWhen operator works similarly to retryWhen, but reacts to successful onComplete events instead of errors.
Example with 3-second retries on failure and 5-second delays after success:
getMyTask() .retryWhen(throwables -> throwables.flatMap(throwable -> { if (throwable instanceof IllegalArgumentException) { Log.d("ERSEN", "Retrying after error..."); return Observable.timer(3, TimeUnit.SECONDS); } return Observable.error(throwable); })) .repeatWhen(completed -> completed.flatMap(ignored -> { Log.d("ERSEN", "Task succeeded, waiting to poll again..."); return Observable.timer(5, TimeUnit.SECONDS); })) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe( result -> Log.d("ERSEN", "Poll result: " + result), error -> Log.d("ERSEN", "Polling stopped with error: " + error.getMessage()) );
Handling Single-Subscriber Requirements
If you need to ensure only one subscriber is running the polling flow (to avoid duplicate concurrent tasks), here are two solid approaches:
Option 1: Cached Shared Observable
Create a single instance of your polling Observable and use share() to let multiple downstream subscribers share the same upstream flow (but only one upstream subscription runs at a time):
private Observable<String> pollingObservable; private Observable<String> getPollingObservable() { if (pollingObservable == null) { pollingObservable = getMyTask() .retryWhen(throwables -> throwables.flatMap(throwable -> { if (throwable instanceof IllegalArgumentException) { return Observable.timer(3, TimeUnit.SECONDS); } return Observable.error(throwable); })) .repeatWhen(completed -> completed.flatMap(ignored -> Observable.timer(5, TimeUnit.SECONDS))) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .share(); } return pollingObservable; } // Usage: Disposable disposable = getPollingObservable().subscribe(...);
Option 2: Controlled Processor with Subscription Check
Use a PublishProcessor to encapsulate the polling flow, and track whether an upstream subscription is already active to prevent duplicates:
private PublishProcessor<String> pollingProcessor = PublishProcessor.create(); private boolean isPollingActive = false; public void startPolling() { if (!isPollingActive) { isPollingActive = true; getMyTask() .retryWhen(throwables -> throwables.flatMap(throwable -> { if (throwable instanceof IllegalArgumentException) { return Observable.timer(3, TimeUnit.SECONDS); } return Observable.error(throwable); })) .repeatWhen(completed -> completed.flatMap(ignored -> Observable.timer(5, TimeUnit.SECONDS))) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe( pollingProcessor::onNext, error -> { pollingProcessor.onError(error); isPollingActive = false; }, () -> { pollingProcessor.onComplete(); isPollingActive = false; } ); } } // External subscription: Disposable disposable = pollingProcessor.subscribe( result -> Log.d("ERSEN", "Poll result: " + result), error -> Log.d("ERSEN", "Polling failed: " + error.getMessage()) );
Don't Forget to Clean Up!
Always dispose of your Disposable when you're done (e.g., in onDestroy for Android Activities/Fragments) to prevent memory leaks:
@Override protected void onDestroy() { super.onDestroy(); if (disposable != null && !disposable.isDisposed()) { disposable.dispose(); } }
内容的提问来源于stack exchange,提问作者Ersen Osman

