从RxJava迁移至RxJava2后retryWhen执行挂起问题排查
RxJava1 到 RxJava2 迁移中 retryWhen 逻辑的问题分析与修复
问题根源
你的代码在RxJava2中出现测试挂起,核心原因是RxJava1与RxJava2对retryWhen操作符的错误流终止逻辑处理存在差异,尤其是在不需要重试时的错误传播方式。
原代码中,当遇到非1/2状态码的ClientException时,你使用了:
errors.flatMap(Observable::error)
这在RxJava1中可能会直接触发流终止,但在RxJava2中,retryWhen的上游错误流errors是冷Observable,再次订阅它会导致没有新的错误事件发出,整个流陷入等待状态,最终导致单元测试挂起。
正确实现方式
我们需要在不需要重试时,直接抛出当前错误来终止retryWhen流,而非复用原错误流。修改后的shouldRetry方法如下:
private Observable<Long> shouldRetry(Observable<? extends Throwable> errors) { return errors.filter(e -> e instanceof ClientException) .cast(ClientException.class) .flatMap(clientEx -> { int statusCode = clientEx.getStatusCode(); if (statusCode == 1 || statusCode == 2) { return Observable.just(statusCode); } else { // 抛出原始异常,终止重试逻辑并传递错误 return Observable.error(clientEx); } }) .zipWith(Observable.range(1, 3), (n, i) -> i) .doOnNext(retryCount -> log.info("Retrying {}", retryCount)) .flatMap(retryCount -> Observable.timer(500, TimeUnit.MILLISECONDS)); }
关键细节说明
RxJava2中retryWhen的错误流行为规则:
- 若错误流发送
onError信号:会终止重试逻辑,并将该错误传递给下游订阅者 - 若错误流发送
onComplete信号:会让原Observable直接完成,不会触发重试 - 若错误流无任何信号:会导致整个订阅挂起,就是你遇到的测试卡住问题
RxJava1与RxJava2核心差异总结
- RxJava1允许通过复用原错误流触发终止,依赖隐式的流终止逻辑;RxJava2则要求明确发送
onError或onComplete信号,冷Observable重复订阅会导致无事件输出 - RxJava2对
retryWhen的错误流信号语义做了严格定义,必须显式控制流的终止行为
内容的提问来源于stack exchange,提问作者Ben Green
相关产品推荐
相关产品推荐

