解决高频消息使用 redis 缓存引起的前端卡死的 bug,高频消息不缓存
This commit is contained in:
parent
c7e16eed9c
commit
2cedd90ae7
@ -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:
|
||||
|
||||
@ -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
|
||||
|
||||
@ -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;
|
||||
|
||||
@ -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
|
||||
|
||||
@ -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<WebSocketSession> 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<WebSocketSession> 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);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Loading…
Reference in New Issue
Block a user