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

从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 18:45:03