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

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.

Core Idea for Data Polling

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 03:40:52