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

如何在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的逻辑
}

最优解决方案

核心思路

  1. 仅针对Token过期类错误触发重试,避免无效重试
  2. 重试前自动刷新Token,全程保持响应式流不阻塞
  3. 优先使用未过期的Token,减少不必要的认证请求
  4. 处理并发场景,避免多个请求重复刷新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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 23:45:40