如何用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线程执行 }
关键细节解释
为什么用
defer?defer会在每次订阅时重新创建内部Observable,这样每次重试都会重新发起gRPC请求,不会复用之前的错误结果。zipWith(Observable.range(1,2))的作用?
用来控制重试次数,这里设置最多重试1次(range从1到2,第一次重试count=1,第二次就超出限制),你可以根据需求调整range的结束值。错误区分的必要性
我们只针对TokenExpiredException执行刷新重试,其他错误(如网络异常、其他业务错误)直接终止,避免无效重试。线程调度
刷新token的操作必须放在IO线程,所以refreshToken()中添加了subscribeOn(Schedulers.io()),避免阻塞主线程。
这样整个流程就非常清晰:请求因token过期失败→触发retryWhen→刷新token→更新全局token→重试请求→成功则处理响应,失败则终止。
内容的提问来源于stack exchange,提问作者Noname135

