如何基于Spring Boot 2与Webflux实现ReactorNettyWebSocketClient重连?
嘿,这个问题我熟!基于Spring Boot 2 + Webflux的ReactorNettyWebSocketClient重连,核心就是利用Reactor的重试机制来自动恢复断开的连接,我给你一步步拆解实现方案:
核心思路
咱们把WebSocket连接的整个生命周期包装成一个Mono,当连接因为超时(你提到的3秒无输入)、服务主动关闭、网络异常等原因断开时,通过Reactor的retryWhen操作符触发重连逻辑。同时可以配置指数退避+随机抖动的重试策略,避免服务恢复时被大量重连请求压垮。
完整代码实现
1. 定义WebSocket客户端服务类
这个类负责创建连接、处理消息、实现重连逻辑:
import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import reactor.core.publisher.Mono; import reactor.netty.http.client.HttpClient; import reactor.netty.http.websocket.WebSocketHandshakeException; import reactor.netty.http.client.ReactorNettyWebSocketClient; import java.net.URI; import java.time.Duration; import java.io.IOException; import java.util.concurrent.TimeoutException; import reactor.util.retry.Retry; @Slf4j @Service public class ReconnectingWebSocketClient { private final ReactorNettyWebSocketClient webSocketClient; private final String targetWsUrl = "ws://your-service-host:port/path"; // 替换成你的服务地址 public ReconnectingWebSocketClient() { // 配置HttpClient,设置3秒无输入超时(对应你提到的超时场景) HttpClient httpClient = HttpClient.create() .responseTimeout(Duration.ofSeconds(3)); this.webSocketClient = new ReactorNettyWebSocketClient(httpClient); } // 启动WebSocket连接(包含重连逻辑) public void startPersistentConnection() { establishConnection() .retryWhen(Retry.backoff(Integer.MAX_VALUE, Duration.ofSeconds(5)) // 加入随机抖动,避免所有客户端同时重连 .jitter(0.5) // 只在指定异常时触发重连 .filter(throwable -> { boolean shouldRetry = throwable instanceof TimeoutException || throwable instanceof IOException || throwable instanceof WebSocketHandshakeException; if (!shouldRetry) { log.error("无需重连的异常", throwable); } return shouldRetry; }) // 重试前打印日志,方便排查 .doBeforeRetry(retrySignal -> log.warn("WebSocket连接断开,准备第{}次重连...", retrySignal.totalRetries() + 1)) ) .subscribe( null, error -> log.error("WebSocket重连耗尽,无法恢复连接", error) ); } // 建立单次WebSocket连接的逻辑 private Mono<Void> establishConnection() { return webSocketClient.execute( URI.create(targetWsUrl), session -> { // 处理收到的消息 Mono<Void> receiveHandler = session.receive() .doOnNext(message -> { // 这里替换成你的消息处理逻辑 log.info("收到WebSocket消息: {}", message.getPayloadAsText()); message.release(); // 记得释放消息资源 }) .doOnError(error -> log.error("消息处理出错", error)) .then(); // 如果需要发送消息,可以在这里添加send逻辑 // Mono<Void> sendHandler = session.send(Mono.just(session.textMessage("heartbeat"))) // 合并接收和发送逻辑,保持连接活跃 return Mono.zip(receiveHandler, Mono.empty()) // 如果有sendHandler就替换Mono.empty() .then(); } ); } }
2. 启动时自动建立连接
用ApplicationRunner在项目启动时触发WebSocket连接:
import lombok.extern.slf4j.Slf4j; import org.springframework.boot.ApplicationArguments; import org.springframework.boot.ApplicationRunner; import org.springframework.stereotype.Component; @Slf4j @Component public class WebSocketStartupTrigger implements ApplicationRunner { private final ReconnectingWebSocketClient webSocketClient; public WebSocketStartupTrigger(ReconnectingWebSocketClient webSocketClient) { this.webSocketClient = webSocketClient; } @Override public void run(ApplicationArguments args) { webSocketClient.startPersistentConnection(); log.info("WebSocket客户端已启动,将自动维护连接"); } }
关键细节解析
重试策略选型:
- 用
Retry.backoff实现指数退避重试,初始间隔5秒,后续间隔会指数增长(可通过maxBackoff限制最大值) jitter(0.5)加入±50%的随机抖动,避免多个客户端在同一时间点重连,减轻服务端压力Integer.MAX_VALUE设置无限重试,你也可以改成固定次数(比如10次)
- 用
异常过滤逻辑:
- 只对
TimeoutException(3秒无输入)、IOException(网络异常)、WebSocketHandshakeException(握手失败)触发重连 - 排除业务逻辑异常,避免因消息解析错误等非连接问题导致不必要的重连
- 只对
连接生命周期维护:
session.receive().then()会让Mono一直处于活跃状态,直到连接主动断开或出错- 如果需要发送消息,只需把发送逻辑的Mono和接收逻辑合并,就能保持连接持续活跃
资源释放:
- 处理消息时记得调用
message.release(),避免ReactorNetty的内存泄漏问题
- 处理消息时记得调用
内容的提问来源于stack exchange,提问作者crixx
相关产品推荐
相关产品推荐

