如何在Mono重试中使用动态参数?处理过期Token刷新重试
响应式WebClient中Token过期自动刷新重试的最优实现
问题描述
使用响应式WebClient调用API时,每个请求需携带可能过期的access_token。需要捕获Token过期导致的请求失败,刷新Token后重试请求,且Token获取通过同一WebClient以Mono形式实现,不能阻塞。
初始尝试代码
private AtomicReference<Token> token; ... public Mono<ApiResponse> callApi() { return Mono.justOrEmpty(token.get()) .switchIfEmpty( Mono.defer(() -> auth() .doOnNext(token::set))) .flatMap(token -> performRequest(token)) .doOnError(e -> { var newToken = auth().block(); // 此处不能阻塞 token.set(newToken); }) .retry(5); } private Mono<Token> auth() { // 调用认证接口返回Token的逻辑 }
已实现的方案(不够优雅)
private AtomicReference<TokenHolder> tokenHolder = new AtomicReference<>(); private AtomicBoolean lastQueryFailed = new AtomicBoolean(); public Mono<ApiResponse> getApiResponse() { return Mono.defer(this::getToken) .flatMap(this::requestApi) .doOnError((e) -> { log.error(e); lastQueryFailed.set(true); }) .retryWhen(Retry.backoff( 3, Duration.ofSeconds(2) )); } private Mono<Token> getToken() { if (tokenHolder.get() == null) { return auth() .doOnNext(tokenHolder::set); } if (!lastQueryFailed.get()) { return Mono.just(tokenHolder.get()); } return auth() .doOnNext(tokenHolder::set); } private Mono<Token> auth() { // 调用认证接口返回Token的逻辑 }
最优解决方案
核心思路
- 仅针对Token过期类错误触发重试,避免无效重试
- 重试前自动刷新Token,全程保持响应式流不阻塞
- 优先使用未过期的Token,减少不必要的认证请求
- 处理并发场景,避免多个请求重复刷新Token
代码实现
import reactor.core.publisher.Mono; import reactor.util.retry.Retry; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.web.reactive.function.client.WebClientResponseException; import java.time.Duration; import java.time.Instant; import java.util.concurrent.atomic.AtomicReference; public class ApiClient { private final WebClient webClient; private final AtomicReference<Token> tokenRef = new AtomicReference<>(); // 并发场景下缓存刷新Token的请求,避免重复调用认证接口 private Mono<Token> cachedRefreshToken; public ApiClient(WebClient webClient) { this.webClient = webClient; } public Mono<ApiResponse> callApi() { return getValidToken() .flatMap(this::performRequest) .retryWhen(Retry.backoff(3, Duration.ofSeconds(2)) // 仅处理Token过期导致的未授权错误 .filter(this::isTokenExpiredError) // 重试前触发Token刷新 .doBeforeRetry(retrySignal -> refreshToken().subscribe(tokenRef::set))); } // 获取有效Token:优先用未过期的缓存Token,否则刷新 private Mono<Token> getValidToken() { Token currentToken = tokenRef.get(); if (currentToken != null && !isTokenExpired(currentToken)) { return Mono.just(currentToken); } return refreshToken() .doOnNext(tokenRef::set); } // 刷新Token,并发场景下缓存请求避免重复调用 private Mono<Token> refreshToken() { if (cachedRefreshToken == null || cachedRefreshToken.isDisposed()) { cachedRefreshToken = webClient.post() .uri("/auth/token") .retrieve() .bodyToMono(Token.class) // 提前30秒失效缓存,避免使用即将过期的Token .cache(token -> Duration.between(Instant.now(), token.getExpiresAt().minusSeconds(30)), error -> Duration.ZERO, () -> Duration.ZERO); } return cachedRefreshToken; } // 判断是否为Token过期导致的错误(根据实际业务调整) private boolean isTokenExpiredError(Throwable error) { return error instanceof WebClientResponseException.Unauthorized; } // 判断Token是否已过期 private boolean isTokenExpired(Token token) { return Instant.now().isAfter(token.getExpiresAt()); } // 执行实际API请求 private Mono<ApiResponse> performRequest(Token token) { return webClient.get() .uri("/api/resource") .header("Authorization", "Bearer " + token.getAccessToken()) .retrieve() .bodyToMono(ApiResponse.class); } // 假设的Token实体类 private static class Token { private String accessToken; private Instant expiresAt; public String getAccessToken() { return accessToken; } public Instant getExpiresAt() { return expiresAt; } } // 假设的API响应实体类 private static class ApiResponse {} }
方案优势
- 精准重试:仅对Token过期的未授权错误重试,避免无关错误触发无效重试
- 无阻塞流:全程使用响应式操作,无
block()调用,符合Reactor编程模型 - 高效缓存:
cachedRefreshToken保证并发请求共享同一个刷新Token的请求,减少认证接口调用量 - 提前失效:缓存Token时提前30秒失效,避免使用即将过期的Token导致请求失败
内容的提问来源于stack exchange,提问作者Seventh
相关产品推荐
相关产品推荐

