Spring Integration中如何检测WebSocket断开并实现重连?
Spring Integration WebSocket 断开检测与自动重连方案
核心思路
通过ClientWebSocketContainer的连接状态监听捕获断开事件,结合自定义重试逻辑或Spring Retry实现自动重连。
具体实现步骤
1. 添加连接状态监听器
给ClientWebSocketContainer绑定WebSocketConnectionListener,监听连接关闭、错误事件:
@Bean IntegrationFlow integrationFlow() { StandardWebSocketClient client = new StandardWebSocketClient(); ClientWebSocketContainer container = new ClientWebSocketContainer(client, url); // 监听连接状态变化 container.addWebSocketListener(new WebSocketConnectionListener() { @Override public void onOpen(WebSocketSession session) { System.out.println("WebSocket连接已建立"); } @Override public void onClose(WebSocketSession session, CloseStatus closeStatus) { System.out.println("WebSocket连接断开,状态码:" + closeStatus.getCode()); // 触发重连 reconnect(container); } @Override public void onError(WebSocketSession session, Throwable exception) { System.err.println("WebSocket连接异常:" + exception.getMessage()); // 异常时触发重连 reconnect(container); } }); WebSocketInboundChannelAdapter adapter = new WebSocketInboundChannelAdapter(container); return IntegrationFlow.from(adapter).handle(m -> System.out.println("收到消息:" + m.getPayload())).get(); }
2. 实现基础重连逻辑
重连时先停止旧容器,延迟后重启建立新连接:
private void reconnect(ClientWebSocketContainer container) { try { // 关闭当前容器 container.stop(); // 3秒后重试,避免频繁请求 Thread.sleep(3000); // 重启容器重建连接 container.start(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.err.println("重连线程被中断"); } catch (Exception e) { System.err.println("重连失败,将再次尝试:" + e.getMessage()); // 递归重试,可按需添加次数限制 reconnect(container); } }
3. 进阶:用Spring Retry控制重试策略
如果需要指数退避、重试次数限制等灵活策略,可引入Spring Retry:
第一步:添加依赖(Maven)
<dependency> <groupId>org.springframework.retry</groupId> <artifactId>spring-retry</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-aop</artifactId> </dependency>
第二步:配置重试模板
@Bean RetryTemplate retryTemplate() { // 最多重试5次 SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); retryPolicy.setMaxAttempts(5); // 指数退避:初始延迟1秒,每次翻倍,最大延迟10秒 ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy(); backOffPolicy.setInitialInterval(1000); backOffPolicy.setMultiplier(2); backOffPolicy.setMaxInterval(10000); RetryTemplate retryTemplate = new RetryTemplate(); retryTemplate.setRetryPolicy(retryPolicy); retryTemplate.setBackOffPolicy(backOffPolicy); return retryTemplate; }
第三步:用RetryTemplate重构重连方法
@Autowired private RetryTemplate retryTemplate; private void reconnect(ClientWebSocketContainer container) { retryTemplate.execute(context -> { container.stop(); Thread.sleep(1000); container.start(); return null; }); }
注意事项
- 确保
ClientWebSocketContainer为单例Bean,避免重复创建实例。 - 生产环境建议替换
System.out为日志框架(如SLF4J),便于问题排查。 - 重连延迟和重试次数需根据服务端限流策略调整,避免被拒绝连接。
内容的提问来源于stack exchange,提问作者rcirne
相关产品推荐
相关产品推荐

