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

如何用RxJava的retryWhen实现GRPC请求错误恢复与重试?

这问题我熟,用RxJava的retryWhen确实是处理这类"错误后重试(先执行修复操作)"场景的最佳实践。不过你现在的代码有两个小问题需要先调整,再加上优雅的retryWhen逻辑:

第一步:把业务错误(错误码非0)转换成RxJava的Error事件

你现在是在subscribe里检查错误码,但retryWhen只能响应Observable的onError事件,所以得把业务错误(比如token过期对应的错误码)转成自定义异常,让上游Observable抛出这个异常,这样retryWhen才能捕获到。

先定义一个自定义异常用于标记token过期场景:

public class TokenExpiredException extends RuntimeException {
    public TokenExpiredException(String message) {
        super(message);
    }
}

然后改造你的请求Observable,去掉cache()(因为cache会缓存第一次结果,重试时不会重新发起请求,不符合我们的需求),换成defer确保每次重试都会重新执行请求逻辑:

Observable<Grpc.MyResponse> requestObservable = Observable.defer(() -> 
    Observable.fromCallable(() -> {
        Grpc.MyRequest request = Grpc.MyRequest.newBuilder()
            .setToken(mToken)
            .build();
        return mStub.mytest(request);
    })
    .map(response -> {
        // 检查业务错误码,把特定错误转成异常
        if (response.getCode() == 401) { // 假设401是token过期的错误码
            throw new TokenExpiredException("Token expired, need refresh");
        } else if (response.getCode() != 0) {
            // 其他业务错误抛出对应异常,终止流程
            throw new RuntimeException("Business error: " + response.getCode());
        }
        return response;
    })
);

第二步:实现带token刷新的retryWhen逻辑

retryWhen的核心是接收一个Observable<Throwable>,然后返回新的Observable:如果返回的Observable发射数据,就触发重试;如果发射错误,就终止序列。这里我们要实现:

  • 只在遇到TokenExpiredException时触发刷新token
  • 限制重试次数(避免无限循环)
  • 刷新成功后更新全局token再重试,失败则直接终止

代码示例:

requestObservable
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .retryWhen(throwables -> throwables
        // 用zipWith限制重试次数,这里设置最多重试1次
        .zipWith(Observable.range(1, 2), (throwable, retryCount) -> 
            new Pair<>(throwable, retryCount)
        )
        .flatMap(pair -> {
            Throwable throwable = pair.first;
            int retryCount = pair.second;

            // 仅处理token过期异常且未超过重试次数的情况
            if (throwable instanceof TokenExpiredException && retryCount <= 1) {
                // 执行刷新token操作,返回Observable<String>发射新token
                return refreshToken()
                    .doOnNext(newToken -> {
                        // 更新全局token,确保下一次请求用新值
                        mToken = newToken;
                    })
                    // 刷新成功后发射数据触发重试
                    .map(newToken -> true);
            } else {
                // 其他异常或重试次数耗尽,传递错误终止序列
                return Observable.error(throwable);
            }
        })
    )
    .subscribe(
        response -> {
            // 处理成功响应
        },
        throwable -> {
            // 处理最终错误(如刷新token失败、非token过期错误)
            throwable.printStackTrace();
        }
    );

第三步:实现refreshToken的Observable

这里需要写一个发起刷新token请求的方法,返回Observable封装结果:

private Observable<String> refreshToken() {
    return Observable.fromCallable(() -> {
        Grpc.RefreshTokenRequest request = Grpc.RefreshTokenRequest.newBuilder()
            .setRefreshToken(mRefreshToken) // 假设你持有刷新token
            .build();
        Grpc.RefreshTokenResponse response = mStub.refreshToken(request);
        if (response.getCode() != 0) {
            throw new RuntimeException("Failed to refresh token: " + response.getCode());
        }
        return response.getNewToken();
    })
    .subscribeOn(Schedulers.io()); // 确保刷新操作在IO线程执行
}

关键细节解释

  1. 为什么用defer?
    defer会在每次订阅时重新创建内部Observable,这样每次重试都会重新发起gRPC请求,不会复用之前的错误结果。

  2. zipWith(Observable.range(1,2))的作用?
    用来控制重试次数,这里设置最多重试1次(range从1到2,第一次重试count=1,第二次就超出限制),你可以根据需求调整range的结束值。

  3. 错误区分的必要性
    我们只针对TokenExpiredException执行刷新重试,其他错误(如网络异常、其他业务错误)直接终止,避免无效重试。

  4. 线程调度
    刷新token的操作必须放在IO线程,所以refreshToken()中添加了subscribeOn(Schedulers.io()),避免阻塞主线程。

这样整个流程就非常清晰:请求因token过期失败→触发retryWhen→刷新token→更新全局token→重试请求→成功则处理响应,失败则终止。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:59:38