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

如何在Flux WebClient中动态更新授权头实现SSE客户端Token刷新

响应式SSE客户端动态Token刷新实现方案

问题根因

当前代码的核心缺陷是:authService.getIdToken()仅在流首次组装、第一次发起请求时被调用,后续retryWhen触发重连时,会直接复用已经构造完成的请求对象,不会重新执行请求头设置逻辑,始终使用第一次生成的旧Token。等Token1小时有效期过了之后,重连时携带的还是过期Token,自然会认证失败。

核心修复思路

必须保证每次发起请求(包括重试触发的重连请求)时,都实时获取当前最新的有效Token,不能在流初始化阶段就把Token值固定下来。

实现代码

基础修复版(适配现有同步AuthService)

用Mono.defer()包裹整个WebClient请求构造逻辑,defer的特性是每次被订阅时(包括重试触发的重新订阅)才会执行内部的请求构造代码,从而每次都能拿到最新Token:

Flux<ServerSideEvent> sseStream = Mono.defer(() -> {
    // 每次重连都会实时执行这行代码,获取最新Token
    String currentToken = authService.getIdToken();
    return webClient.get()
            .uri("/events")
            .headers(headers -> headers.setBearerAuth(currentToken))
            .retrieve()
            .bodyToFlux(ServerSideEvent.class);
})
.timeout(Duration.ofSeconds(TIMEOUT))
.retryWhen(Retry.fixedDelay(Long.MAX_VALUE, Duration.ofSeconds(RETRY_DELAY))
        .doBeforeRetry(signal -> {
            // 如果是401未授权错误,先强制刷新Token再重试
            if (signal.failure() instanceof WebClientResponseException.Unauthorized) {
                authService.refreshToken();
            }
        })
)
.subscribe(
        this::handleEvent,
        err -> logger.error("SSE连接异常: {}", err.getMessage()),
        () -> logger.info("SSE连接断开")
);

无阻塞响应式版(适配响应式AuthService)

如果你的AuthService本身是响应式实现,不要在响应式流里调用阻塞方法,直接把Token获取逻辑纳入流编排:

Flux<ServerSideEvent> sseStream = Mono.defer(() ->
    // getValidToken() 内部实现Token缓存、过期自动刷新逻辑,返回Mono<String>
    authService.getValidToken()
        .flatMapMany(validToken -> webClient.get()
                .uri("/events")
                .headers(headers -> headers.setBearerAuth(validToken))
                .retrieve()
                .bodyToFlux(ServerSideEvent.class)
        )
)
.timeout(Duration.ofSeconds(TIMEOUT))
// 推荐用退避重试替代固定间隔重试,避免服务端故障时请求压力过大
.retryWhen(Retry.backoff(Long.MAX_VALUE, Duration.ofSeconds(1))
        .maxBackoff(Duration.ofSeconds(RETRY_DELAY))
        // 过滤不可重试的错误,比如403权限错误直接终止流
        .filter(err -> !(err instanceof WebClientResponseException.Forbidden))
)
.subscribe(
        this::handleEvent,
        err -> logger.error("SSE连接异常: {}", err.getMessage()),
        () -> logger.info("SSE连接断开")
);

优化建议

  • AuthService层提前做Token有效期校验:比如检测到Token剩余有效期不足5分钟时主动刷新,不要等收到401或者连接断开才刷新,降低重连失败概率
  • 增加心跳检测逻辑:如果超过指定时长没有收到任何SSE事件(包括服务端心跳),主动断开连接触发重连,避免连接假死
  • 不要在全局共享的WebClient实例上提前固定Authorization头,必须在每次构造单个请求时动态设置最新Token
  • 重试时增加日志埋点,记录重连原因、重连次数,方便排查问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 23:06:27