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

Spring @ClientEndpoint实现的WebSocket客户端无法检测网线拔除的网络断开问题

问题原因分析
  • 你未实现pong消息的接收和超时判断逻辑:当前代码的@OnMessage仅处理文本类型消息,无法接收服务端返回的Pong帧,你单向发送ping但从未验证是否收到响应,即使网络断开,只要本地会话标记未变更,就会一直认为连接正常。
  • Session.isOpen()仅校验本地会话状态:该方法不会真实探测网络连通性,拔除网线属于静默断网,没有TCP FIN/RST包交互,本地TCP栈和WebSocket实现都无法感知连接中断,会一直标记会话为打开状态。
  • 默认无空闲超时配置:JSR 356(Java WebSocket API)默认未开启会话空闲超时,只会在网络恢复后,旧连接的报文被服务端重置时,才会触发1006异常,也就是你看到的Connection reset by peer报错。
修复方案

1. 新增Pong消息处理与心跳超时检测

修改WSClient类,新增最后活跃时间统计、Pong帧处理逻辑,定时判断心跳超时:

import javax.websocket.*;
import java.io.IOException;
import java.nio.ByteBuffer;

@ClientEndpoint
public class WSClient {
    private Session session;
    private int i = 0;
    // 上次收到服务端消息的时间戳
    private volatile long lastActiveTime;
    // 心跳超时阈值10秒,可自行调整
    private static final long HEARTBEAT_TIMEOUT = 10 * 1000;

    @OnOpen
    public void open(Session session) {
        System.out.println("Connected to the server");
        this.session = session;
        this.lastActiveTime = System.currentTimeMillis();
    }

    @OnClose
    public void close(Session session, CloseReason closeReason) {
        System.out.println("connection closed " + closeReason.getReasonPhrase());
    }

    @OnError
    public void error(Session session, Throwable t) {
        System.out.println(session.getId());
        System.out.println("Error in connection " + t.getMessage());
    }

    // 处理文本消息
    @OnMessage
    public void message(Session session, String message) {
        this.lastActiveTime = System.currentTimeMillis();
        System.out.println("message received: " + message + " " + i++);
    }

    // 新增:处理Pong消息
    @OnMessage
    public void handlePong(Session session, PongMessage pong) {
        this.lastActiveTime = System.currentTimeMillis();
        System.out.println("received pong from server");
    }

    public void send(String message){
        try {
            // 先判断是否心跳超时
            if (System.currentTimeMillis() - lastActiveTime > HEARTBEAT_TIMEOUT) {
                System.out.println("heartbeat timeout, close connection");
                session.close(new CloseReason(CloseReason.CloseCodes.GOING_AWAY, "heartbeat timeout"));
                return;
            }
            if(session.isOpen()) {
                this.session.getBasicRemote().sendPing(ByteBuffer.wrap(message.getBytes()));
                System.out.println("send ping success " + i++);
            } else {
                System.out.println("socket closed");
            }
        } catch (IOException e) {
            e.printStackTrace();
            // 发送异常直接标记连接失效
            try {
                session.close();
            } catch (IOException ex) {
                // 忽略关闭异常
            }
        }
    }

    public boolean isAlive() {
        return session != null && session.isOpen() && System.currentTimeMillis() - lastActiveTime <= HEARTBEAT_TIMEOUT;
    }
}

2. 配置会话超时与自动重连逻辑

修改ClientApp类,添加容器超时配置,以及连接断开后的自动重连逻辑,同时调整线程模型避免阻塞Spring启动流程:

import org.springframework.context.ApplicationListener;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.stereotype.Component;
import javax.websocket.ContainerProvider;
import javax.websocket.WebSocketContainer;
import java.net.URI;
import java.util.concurrent.TimeUnit;

@Component
public class ClientApp implements ApplicationListener<ApplicationReadyEvent> {

    private WSClient client;
    // 重连间隔5秒
    private static final long RECONNECT_INTERVAL = 5 * 1000;

    private void startConnection() throws Exception {
        WebSocketContainer container = ContainerProvider.getWebSocketContainer();
        // 配置会话最大空闲时间15秒
        container.setDefaultMaxSessionIdleTimeout(15 * 1000);
        client = new WSClient();
        container.connectToServer(client, new URI("ws://wshost:8080/ping"));

        while (true) {
            if (client.isAlive()) {
                client.send("ping");
            } else {
                System.out.println("connection is down, try to reconnect");
                // 重新建立连接
                try {
                    container.connectToServer(client, new URI("ws://wshost:8080/ping"));
                } catch (Exception e) {
                    System.out.println("reconnect failed: " + e.getMessage());
                    Thread.sleep(RECONNECT_INTERVAL);
                    continue;
                }
            }
            TimeUnit.SECONDS.sleep(3);
        }
    }

    @Override
    public void onApplicationEvent(ApplicationReadyEvent event) {
        // 启动单独线程跑WebSocket逻辑,避免阻塞Spring启动流程
        new Thread(() -> {
            while (true) {
                try {
                    startConnection();
                } catch (Exception e) {
                    System.out.println("connection error: " + e.getMessage());
                    try {
                        Thread.sleep(RECONNECT_INTERVAL);
                    } catch (InterruptedException ex) {
                        Thread.currentThread().interrupt();
                        break;
                    }
                }
            }
        }).start();
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 00:24:01