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

