From f104f2f6ac56fcee17e5f505d845273fb5a1d6f5 Mon Sep 17 00:00:00 2001 From: Tian jianyong <11429339@qq.com> Date: Fri, 21 Nov 2025 13:31:31 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8DADXP=E9=80=82=E9=85=8D?= =?UTF-8?q?=E5=99=A8=E7=BC=96=E8=AF=91=E9=94=99=E8=AF=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../service/MessageListenerService.java | 47 +++++++++++-------- 1 file changed, 27 insertions(+), 20 deletions(-) diff --git a/adxp-adapter/src/main/java/com/qaup/adxp/adapter/service/MessageListenerService.java b/adxp-adapter/src/main/java/com/qaup/adxp/adapter/service/MessageListenerService.java index 40d9729a..1ea333f6 100644 --- a/adxp-adapter/src/main/java/com/qaup/adxp/adapter/service/MessageListenerService.java +++ b/adxp-adapter/src/main/java/com/qaup/adxp/adapter/service/MessageListenerService.java @@ -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 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 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()); } } \ No newline at end of file