RSocket-Java流场景下,如何在Flux.retryWhen中更新JWT令牌?
解决方案:RSocket重试时更新JWT令牌
问题原因
当前代码在初始化Flux时就调用了getAccessToken()获取JWT令牌,重试逻辑只是重新订阅已创建的Flux,并不会重新执行令牌获取逻辑,导致重试时使用的是初始的过期令牌,触发RejectedSetupException错误。
解决思路
使用Flux.defer()延迟请求逻辑的执行,让每次重试都重新执行令牌获取和请求发送的完整流程,确保每次重试都使用最新的有效JWT令牌。
修改后的客户端流订阅代码
public void stream(final MessageCallback callback, final StreamRequest streamRequest) { final Flux<GenericRecord> streamFlux = Flux.defer(() -> this.rSocketRequester .route("stream") .metadata(getAccessToken(), AUTHENTICATION) .data(streamRequest) .retrieveFlux(GenericRecord.class) ).retryWhen(retryBackoffSpec); streamFlux.publishOn(scheduler).subscribe(m -> callback.doOnMessage(m), t -> callback.doOnError(t)); }
原理说明
Flux.defer()会在每次订阅(包括重试触发的重新订阅)时才执行内部的Lambda表达式,这样每次重试都会重新调用getAccessToken()获取最新的JWT令牌,再将其作为元数据发送给RSocket服务端,避免了过期令牌的问题。
额外注意事项
- 确保
getAccessToken()方法内部已实现令牌过期检测与刷新逻辑,能返回有效的新令牌。 - 若RSocket连接本身需要重连(而非流订阅重试),客户端初始化时配置的
connector.reconnect(retryBackoffSpec)会处理连接层面的重连,但流订阅的重试仍需通过defer确保令牌更新。
内容的提问来源于stack exchange,提问作者Maciej Lach
相关产品推荐
相关产品推荐

