修复ADXP适配器编译错误

This commit is contained in:
Tian jianyong 2025-11-21 13:31:31 +08:00
parent f08a4dcf0a
commit f104f2f6ac

View File

@ -56,31 +56,36 @@ public class MessageListenerService {
// 获取所有活跃的会话ID
while (isRunning.get()) {
try {
// 获取当前会话数量
int sessionCount = adxpSdkService.getSessions().size();
// 检查ADXP连接状态
if (!adxpSdkService.isConnected()) {
log.debug("ADXP数据中台连接未建立等待重连...");
Thread.sleep(2000);
continue;
}
// 获取当前会话数量单连接架构下为1或0
int sessionCount = adxpSdkService.getActiveSessionCount();
if (sessionCount == 0) {
log.debug("当前没有活跃的会话,等待连接...");
Thread.sleep(1000);
continue;
}
// 遍历所有会话并接收消息
adxpSdkService.getSessions().forEach((sessionId, sessionInfo) -> {
try {
// 调用SDK接收消息这可能会阻塞直到有消息到达
java.util.List<com.qaup.adxp.adapter.dto.FlightMessage> messages =
adxpSdkService.receiveMessages(sessionId);
// 如果有消息广播给所有WebSocket客户端
if (messages != null && !messages.isEmpty()) {
adxpWebSocketHandler.broadcastMessages(messages);
log.info("接收到 {} 条航班消息并广播给 {} 个WebSocket客户端",
messages.size(), adxpWebSocketHandler.getSessionCount());
}
} catch (Exception e) {
log.error("处理会话消息失败: sessionId={}", sessionId, e);
// 单连接架构下直接接收消息
try {
// 调用SDK接收消息
java.util.List<com.qaup.adxp.adapter.dto.FlightMessage> messages =
adxpSdkService.getMessages();
// 如果有消息广播给所有WebSocket客户端
if (messages != null && !messages.isEmpty()) {
adxpWebSocketHandler.broadcastMessages(messages);
log.info("接收到 {} 条航班消息并广播给 {} 个WebSocket客户端",
messages.size(), adxpWebSocketHandler.getSessionCount());
}
});
} catch (Exception e) {
log.error("处理ADXP消息失败", e);
}
// 短暂休眠避免过于频繁的轮询
Thread.sleep(100);
@ -110,10 +115,12 @@ public class MessageListenerService {
public String getStats() {
return String.format("MessageListenerService Stats: " +
"isRunning=%s, " +
"sessions=%d, " +
"adxpConnected=%s, " +
"adxpActiveSessions=%d, " +
"webSocketClients=%d",
isRunning.get(),
adxpSdkService.getSessions().size(),
adxpSdkService.isConnected(),
adxpSdkService.getActiveSessionCount(),
adxpWebSocketHandler.getSessionCount());
}
}