如何在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
相关产品推荐
相关产品推荐

