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

TimeoutException抛出时如何取消Flowable上游的阻塞逻辑?

RxJava2超时后终止上游阻塞逻辑的解决方案

问题场景

你提供的RxJava2代码中,当timeout(1000, TimeUnit.MILLISECONDS)触发TimeoutException后,上游的block()方法仍会继续执行,一段时间后打印"block end"日志。需要实现超时后取消上游逻辑,让block()方法停止运行。

原代码

public class Demo {

    /** @noinspection ResultOfMethodCallIgnored*/
    @SuppressLint("CheckResult")
    public static void test() {
        Flowable.just(false).map(aBoolean -> block())
                .retryWhen(throwableFlowable -> {
                    AtomicInteger retryCounter = new AtomicInteger();
                    return throwableFlowable.takeWhile(throwable -> {
                        if (retryCounter.getAndIncrement() < 5) {
                            return true;
                        }

                        Exception exception;
                        try {
                            exception = (Exception) throwable;
                        } catch (Exception e) {
                            exception = new Exception(throwable);
                        }
                        throw exception;
                    });
                }).timeout(1000, TimeUnit.MILLISECONDS).subscribe(aBoolean -> {
                    Log.e("Demo", "onNext");
                }, throwable -> {
                    Log.e("Demo", "onError", throwable);
                }, () -> {
                    Log.e("Demo", "onComplete");
                }, subscription -> {
                    subscription.request(Long.MAX_VALUE);
                });
    }

    private static boolean block() {
        Log.e("Demo", "block");
        try {
            Thread.sleep(1500);
        } catch (InterruptedException e) {
            // do nothing
        }
        Log.e("Demo", "block end");
        return true;
    }

}

可行解决方案

1. 响应线程中断信号

RxJava在订阅取消(如超时触发终止)时,会中断执行上游任务的线程。只需修改block()方法,不再忽略InterruptedException,而是利用这个中断信号终止方法执行:

private static boolean block() {
    Log.e("Demo", "block");
    try {
        Thread.sleep(1500);
    } catch (InterruptedException e) {
        // 响应中断,恢复线程中断状态并终止方法
        Thread.currentThread().interrupt();
        Log.e("Demo", "block interrupted, exiting");
        return false;
    }
    Log.e("Demo", "block end");
    return true;
}

同时确保上游操作符运行在支持中断的调度器上(比如Schedulers.io()),若未指定调度器,任务可能在调用线程执行,此时需结合subscribeOn使用。

2. 结合Disposable主动取消订阅

将阻塞任务放到IO调度器线程,在超时触发onError时,主动调用Disposable.dispose()取消订阅,触发上游线程中断:

@SuppressLint("CheckResult")
public static void test() {
    Disposable disposable = Flowable.just(false)
            .subscribeOn(Schedulers.io()) // 将阻塞任务分配到IO线程
            .map(aBoolean -> block())
            .retryWhen(throwableFlowable -> {
                AtomicInteger retryCounter = new AtomicInteger();
                return throwableFlowable.takeWhile(throwable -> {
                    if (retryCounter.getAndIncrement() < 5) {
                        return true;
                    }
                    Exception exception;
                    try {
                        exception = (Exception) throwable;
                    } catch (Exception e) {
                        exception = new Exception(throwable);
                    }
                    throw exception;
                });
            })
            .timeout(1000, TimeUnit.MILLISECONDS)
            .subscribe(aBoolean -> {
                Log.e("Demo", "onNext");
            }, throwable -> {
                Log.e("Demo", "onError", throwable);
                disposable.dispose(); // 主动取消订阅,终止上游任务
            }, () -> {
                Log.e("Demo", "onComplete");
            }, subscription -> {
                subscription.request(Long.MAX_VALUE);
            });
}

超时后,dispose()会终止IO线程上的block()任务,Thread.sleep()抛出InterruptedException,我们在block()中响应这个中断即可停止方法。

3. 自定义取消信号(针对非睡眠阻塞逻辑)

如果阻塞逻辑不是Thread.sleep(),而是自定义的循环等待,可通过AtomicBoolean作为取消信号,在循环中检查信号状态:

private static boolean block(AtomicBoolean isCancelled) {
    Log.e("Demo", "block");
    long startTime = System.currentTimeMillis();
    while (System.currentTimeMillis() - startTime < 1500) {
        if (isCancelled.get()) { // 检查取消信号
            Log.e("Demo", "block cancelled, exiting");
            return false;
        }
        try {
            Thread.sleep(100); // 短间隔睡眠,方便检测取消信号
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            Log.e("Demo", "block interrupted, exiting");
            return false;
        }
    }
    Log.e("Demo", "block end");
    return true;
}

修改test()方法,创建并管理这个取消信号:

@SuppressLint("CheckResult")
public static void test() {
    AtomicBoolean isCancelled = new AtomicBoolean(false);
    Disposable disposable = Flowable.just(false)
            .subscribeOn(Schedulers.io())
            .map(aBoolean -> block(isCancelled))
            .retryWhen(throwableFlowable -> {
                AtomicInteger retryCounter = new AtomicInteger();
                return throwableFlowable.takeWhile(throwable -> {
                    if (retryCounter.getAndIncrement() < 5) {
                        isCancelled.set(false); // 重试时重置取消信号
                        return true;
                    }
                    Exception exception;
                    try {
                        exception = (Exception) throwable;
                    } catch (Exception e) {
                        exception = new Exception(throwable);
                    }
                    throw exception;
                });
            })
            .timeout(1000, TimeUnit.MILLISECONDS)
            .subscribe(aBoolean -> {
                Log.e("Demo", "onNext");
            }, throwable -> {
                Log.e("Demo", "onError", throwable);
                isCancelled.set(true); // 触发取消信号
                disposable.dispose();
            }, () -> {
                Log.e("Demo", "onComplete");
            }, subscription -> {
                subscription.request(Long.MAX_VALUE);
            });
}

内容的提问来源于stack exchange,提问作者Taku shousei

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 19:04:55