RxJava2 amb操作符抛出java.lang.InterruptedException问题求助
解决RxJava2 ambArray操作符抛出InterruptedException的问题
你遇到的这个异常其实是ambArray操作符的正常行为导致的——当其中一个Observable率先发射数据后,ambArray会立刻dispose掉其余所有Observable,这就会让那些还在执行阻塞操作(比如你代码里的Thread.sleep())的线程被中断,进而抛出InterruptedException。下面给你两个实用的解决思路:
1. 捕获中断异常并检查Observable的废弃状态
既然这种中断是ambArray的预期逻辑,我们可以在代码里捕获InterruptedException,同时通过emitter.isDisposed()判断当前Observable是否已经被废弃:如果是,就直接结束逻辑,不用把异常向上抛出;如果不是,再把异常传递给下游。
修改后的Observable1代码示例:
Observable<String> observable1 = Observable.defer(new Callable<ObservableSource<? extends String>>() { @Override public ObservableSource<? extends String> call() throws Exception { return Observable.create(new ObservableOnSubscribe<String>() { @Override public void subscribe(final ObservableEmitter<String> emitter) throws Exception { try { long sleepTime = getRandomSleepTime(); // 你的随机睡眠时间逻辑 Thread.sleep(sleepTime); // 只有Observable未被废弃时,才发射数据 if (!emitter.isDisposed()) { emitter.onNext("Observable1 完成发射"); emitter.onComplete(); } } catch (InterruptedException e) { // 仅当Observable未被废弃时,才传递错误 if (!emitter.isDisposed()) { emitter.onError(e); } // 重置线程中断状态(可选,根据你的线程上下文决定) Thread.currentThread().interrupt(); } } }); } });
2. 用RxJava内置的delay操作符替代Thread.sleep()
手动调用Thread.sleep()并不是RxJava的推荐写法,我们可以用delay操作符实现延迟发射,它会自动处理线程中断的情况,不需要你手动捕获异常,代码也更简洁:
优化后的Observable1代码示例:
Observable<String> observable1 = Observable.defer(() -> Observable.just("Observable1 完成发射") .delay(getRandomSleepTime(), TimeUnit.MILLISECONDS) );
这种写法完全遵循RxJava的响应式编程规范,从根源上避免了手动操作线程带来的异常问题。
最后再提一句:ambArray的核心就是“先到先得”,只要有一个Observable发射数据,其余Observable就会被立即取消订阅,所以阻塞操作被中断是正常现象,我们只需要正确处理这种场景就好。
内容的提问来源于stack exchange,提问作者Abu Yousuf
相关产品推荐
相关产品推荐

