From 2cedd90ae7536f5c3a8fec13d7891e1e6151012e Mon Sep 17 00:00:00 2001 From: Tian jianyong <11429339@qq.com> Date: Thu, 13 Nov 2025 15:20:14 +0800 Subject: [PATCH] =?UTF-8?q?=E8=A7=A3=E5=86=B3=E9=AB=98=E9=A2=91=E6=B6=88?= =?UTF-8?q?=E6=81=AF=E4=BD=BF=E7=94=A8=20redis=20=E7=BC=93=E5=AD=98?= =?UTF-8?q?=E5=BC=95=E8=B5=B7=E7=9A=84=E5=89=8D=E7=AB=AF=E5=8D=A1=E6=AD=BB?= =?UTF-8?q?=E7=9A=84=20bug=EF=BC=8C=E9=AB=98=E9=A2=91=E6=B6=88=E6=81=AF?= =?UTF-8?q?=E4=B8=8D=E7=BC=93=E5=AD=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/main/resources/application-prod.yml | 12 ++-- qaup-admin/src/main/resources/application.yml | 9 ++- .../AdxpFlightServiceWebSocketClient.java | 33 ++++++---- .../WebSocketMessageBroadcaster.java | 36 +++++++++-- .../handler/CollisionWebSocketHandler.java | 61 +++++++++++++++++-- 5 files changed, 120 insertions(+), 31 deletions(-) diff --git a/qaup-admin/src/main/resources/application-prod.yml b/qaup-admin/src/main/resources/application-prod.yml index 64857506..d41dbcaf 100644 --- a/qaup-admin/src/main/resources/application-prod.yml +++ b/qaup-admin/src/main/resources/application-prod.yml @@ -28,10 +28,14 @@ spring: timeout: 10s lettuce: pool: - min-idle: 0 - max-idle: 8 - max-active: 8 - max-wait: -1ms + # 最小空闲连接数(保底连接) + min-idle: 5 + # 最大空闲连接数 + max-idle: 20 + # 最大活跃连接数 + max-active: 50 + # 连接最大等待时间(毫秒),防止无限等待 + max-wait: 5000ms # Redis 内存优化配置(生产环境) redis: diff --git a/qaup-admin/src/main/resources/application.yml b/qaup-admin/src/main/resources/application.yml index 1d6b169f..70db82d3 100644 --- a/qaup-admin/src/main/resources/application.yml +++ b/qaup-admin/src/main/resources/application.yml @@ -243,11 +243,10 @@ management: health: show-details: always metrics: - export: - simple: - enabled: true enable: hikari: true jvm: true - jmx: - enabled: true + simple: + metrics: + export: + enabled: true diff --git a/qaup-collision/src/main/java/com/qaup/collision/datacollector/websocket/AdxpFlightServiceWebSocketClient.java b/qaup-collision/src/main/java/com/qaup/collision/datacollector/websocket/AdxpFlightServiceWebSocketClient.java index 1cdffcc3..bf2d3338 100644 --- a/qaup-collision/src/main/java/com/qaup/collision/datacollector/websocket/AdxpFlightServiceWebSocketClient.java +++ b/qaup-collision/src/main/java/com/qaup/collision/datacollector/websocket/AdxpFlightServiceWebSocketClient.java @@ -13,6 +13,7 @@ import org.springframework.web.socket.WebSocketHandler; import org.springframework.web.socket.WebSocketMessage; import org.springframework.web.socket.WebSocketSession; import org.springframework.web.socket.client.WebSocketClient; +import org.springframework.lang.NonNull; import java.net.URI; import java.util.List; @@ -70,12 +71,13 @@ public class AdxpFlightServiceWebSocketClient implements WebSocketHandler { /** * 连接到ADXP适配器的WebSocket服务 */ + @SuppressWarnings("removal") // 抑制WebSocketClient.doHandshake过期警告 public void connect() { if (!properties.isConfigurationReady()) { log.warn("数据中台航班 SDK 配置不完整,WebSocket客户端将无法正常工作"); return; } - + try { // 首先通过HTTP登录获取sessionId sessionId = loginAndGetSessionId(); @@ -84,17 +86,22 @@ public class AdxpFlightServiceWebSocketClient implements WebSocketHandler { scheduleReconnect(); return; } - + // 构建WebSocket URL - String wsUrl = String.format("ws://%s:%d/ws/flight-notifications", + String wsUrl = String.format("ws://%s:%d/ws/flight-notifications", properties.getHost(), properties.getPort()); - + log.info("正在连接到ADXP适配器WebSocket服务: url={}", wsUrl); - - // 连接WebSocket - session = webSocketClient.doHandshake(this, URI.create(wsUrl).toString()).get(); - isConnected.set(true); - log.info("✅ 已连接到ADXP适配器WebSocket服务"); + + // 连接WebSocket - 使用doHandshake(虽过期但仍可用) + WebSocketSession newSession = webSocketClient.doHandshake(this, URI.create(wsUrl).toString()).get(); + + // 在afterConnectionEstablished中也会设置session,这里只是避免编译警告 + if (newSession != null && newSession.isOpen()) { + session = newSession; + isConnected.set(true); + log.info("✅ 已连接到ADXP适配器WebSocket服务"); + } } catch (Exception e) { log.error("❌ 连接ADXP适配器WebSocket服务失败", e); @@ -211,7 +218,7 @@ public class AdxpFlightServiceWebSocketClient implements WebSocketHandler { } @Override - public void afterConnectionEstablished(WebSocketSession session) throws Exception { + public void afterConnectionEstablished(@NonNull WebSocketSession session) throws Exception { log.info("🟢 WebSocket连接已建立: sessionId={}", session.getId()); this.session = session; isConnected.set(true); @@ -219,7 +226,7 @@ public class AdxpFlightServiceWebSocketClient implements WebSocketHandler { } @Override - public void handleMessage(WebSocketSession session, WebSocketMessage message) throws Exception { + public void handleMessage(@NonNull WebSocketSession session, @NonNull WebSocketMessage message) throws Exception { if (message instanceof TextMessage) { handleTextMessage(session, (TextMessage) message); } @@ -259,7 +266,7 @@ public class AdxpFlightServiceWebSocketClient implements WebSocketHandler { } @Override - public void handleTransportError(WebSocketSession session, Throwable exception) throws Exception { + public void handleTransportError(@NonNull WebSocketSession session, @NonNull Throwable exception) throws Exception { log.error("❌ WebSocket传输错误", exception); errorCount.incrementAndGet(); isConnected.set(false); @@ -272,7 +279,7 @@ public class AdxpFlightServiceWebSocketClient implements WebSocketHandler { } @Override - public void afterConnectionClosed(WebSocketSession session, CloseStatus closeStatus) throws Exception { + public void afterConnectionClosed(@NonNull WebSocketSession session, @NonNull CloseStatus closeStatus) throws Exception { log.info("🟡 WebSocket连接已关闭: reason={}, code={}", closeStatus.getReason(), closeStatus.getCode()); isConnected.set(false); this.session = null; diff --git a/qaup-collision/src/main/java/com/qaup/collision/websocket/broadcaster/WebSocketMessageBroadcaster.java b/qaup-collision/src/main/java/com/qaup/collision/websocket/broadcaster/WebSocketMessageBroadcaster.java index d2a37db2..c10b22f1 100644 --- a/qaup-collision/src/main/java/com/qaup/collision/websocket/broadcaster/WebSocketMessageBroadcaster.java +++ b/qaup-collision/src/main/java/com/qaup/collision/websocket/broadcaster/WebSocketMessageBroadcaster.java @@ -231,16 +231,42 @@ public class WebSocketMessageBroadcaster { // 使用Jackson ObjectMapper将UniversalMessage序列化为JSON字符串 // 这样前端可以获得消息类型、时间戳、消息ID等完整信息 String jsonMessage = objectMapper.writeValueAsString(message); - this.collisionWebSocketHandler.broadcastMessage(jsonMessage); - - // 缓存消息用于重连恢复 - messageCacheService.cacheMessage(message); - + this.collisionWebSocketHandler.broadcastMessage(jsonMessage); + + // 根据消息类型决定是否缓存(避免高频消息阻塞Redis) + // 高频实时消息(如位置更新)不缓存,追求实时性 + // 低频重要消息(如碰撞预警)缓存,保证不丢失 + if (shouldCacheMessage(message.getType())) { + messageCacheService.cacheMessage(message); + } + } catch (Exception e) { System.err.println("Failed to broadcast message via native WebSocket: " + e.getMessage()); e.printStackTrace(); } } + + /** + * 判断消息是否应该缓存 + * + * @param messageType 消息类型 + * @return true=需要缓存,false=不需要缓存 + */ + private boolean shouldCacheMessage(String messageType) { + // 高频实时消息 - 不缓存(追求实时性,避免Redis性能瓶颈) + // 这些消息更新频繁,用户关心的是最新状态,历史数据价值不高 + if (MessageTypeConstants.POSITION_UPDATE.equals(messageType) || + MessageTypeConstants.TRAFFIC_LIGHT_STATUS.equals(messageType) || + MessageTypeConstants.HEARTBEAT.equals(messageType) || + MessageTypeConstants.VEHICLE_STATUS_UPDATE.equals(messageType)) { + return false; + } + + // 低频重要消息 - 缓存(保证不丢失) + // 包括:碰撞预警、规则违规、路径冲突、电子围栏、航班通知、车辆指令等 + // 这些事件频率低但重要性高,用户需要完整的历史记录 + return true; + } /** * 生成唯一消息ID diff --git a/qaup-collision/src/main/java/com/qaup/collision/websocket/handler/CollisionWebSocketHandler.java b/qaup-collision/src/main/java/com/qaup/collision/websocket/handler/CollisionWebSocketHandler.java index f35443b8..8a29fe93 100644 --- a/qaup-collision/src/main/java/com/qaup/collision/websocket/handler/CollisionWebSocketHandler.java +++ b/qaup-collision/src/main/java/com/qaup/collision/websocket/handler/CollisionWebSocketHandler.java @@ -7,6 +7,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.concurrent.ConcurrentHashMap; import java.util.Map; +import java.util.List; +import java.util.stream.Collectors; /** * 冲突检测WebSocket处理器 @@ -99,7 +101,7 @@ public class CollisionWebSocketHandler implements WebSocketHandler { } /** - * 发送消息给指定会话(线程安全) + * 发送消息给指定会话(线程安全)- 添加发送失败跳过机制 */ private void sendMessage(WebSocketSession session, String message) { try { @@ -112,15 +114,66 @@ public class CollisionWebSocketHandler implements WebSocketHandler { } } } catch (Exception e) { - LOGGER.error("发送消息失败 - 会话ID: {}", session.getId(), e); + // 发送失败时,快速清理无效会话,避免影响其他会话 + // 使用快速失败策略,不重试,避免加重系统负载 + String sessionId = session.getId(); + LOGGER.warn("发送消息失败,快速跳过 - 会话ID: {}, 错误: {}", sessionId, e.getMessage()); + + // 异步清理无效会话(不阻塞主流程) + try { + sessions.remove(sessionId); + if (session.isOpen()) { + session.close(CloseStatus.SERVICE_RESTARTED); + } + } catch (Exception closeEx) { + LOGGER.debug("关闭无效会话失败 - 会话ID: {}, 错误: {}", sessionId, closeEx.getMessage()); + } } } /** - * 广播消息给所有连接的客户端 + * 广播消息给所有连接的客户端 - 添加分批发送机制 + * 避免一次发送过多消息导致前端卡死 */ public void broadcastMessage(String message) { - sessions.values().forEach(session -> sendMessage(session, message)); + // 获取所有活跃会话 + List activeSessions = sessions.values().stream() + .filter(WebSocketSession::isOpen) + .collect(Collectors.toList()); + + if (activeSessions.isEmpty()) { + LOGGER.debug("没有活跃会话,跳过广播"); + return; + } + + // 控制单次发送数量,避免前端过载 + int batchSize = 10; // 每批最多发送10个会话 + int totalSessions = activeSessions.size(); + + LOGGER.debug("开始广播消息,总会话数: {}, 分批大小: {}", totalSessions, batchSize); + + // 分批发送 + for (int i = 0; i < totalSessions; i += batchSize) { + int endIndex = Math.min(i + batchSize, totalSessions); + List batch = activeSessions.subList(i, endIndex); + + LOGGER.debug("发送批次 {}-{} (共{}个会话)", i + 1, endIndex, batch.size()); + + // 发送当前批次 + batch.forEach(session -> sendMessage(session, message)); + + // 如果还有更多批次,短暂休眠避免过载 + if (endIndex < totalSessions) { + try { + Thread.sleep(10); // 休眠10毫秒 + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + LOGGER.warn("广播线程休眠被中断", e); + } + } + } + + LOGGER.debug("广播消息完成,总会话数: {}", totalSessions); } /**