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

如何基于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客户端已启动,将自动维护连接");
    }
}
关键细节解析
  1. 重试策略选型:

    • 用Retry.backoff实现指数退避重试,初始间隔5秒,后续间隔会指数增长(可通过maxBackoff限制最大值)
    • jitter(0.5)加入±50%的随机抖动,避免多个客户端在同一时间点重连,减轻服务端压力
    • Integer.MAX_VALUE设置无限重试,你也可以改成固定次数(比如10次)
  2. 异常过滤逻辑:

    • 只对TimeoutException(3秒无输入)、IOException(网络异常)、WebSocketHandshakeException(握手失败)触发重连
    • 排除业务逻辑异常,避免因消息解析错误等非连接问题导致不必要的重连
  3. 连接生命周期维护:

    • session.receive().then()会让Mono一直处于活跃状态,直到连接主动断开或出错
    • 如果需要发送消息,只需把发送逻辑的Mono和接收逻辑合并,就能保持连接持续活跃
  4. 资源释放:

    • 处理消息时记得调用message.release(),避免ReactorNetty的内存泄漏问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 03:32:24