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
相关产品推荐
相关产品推荐

