Просмотр исходного кода

feat(websocket): 添加设备状态实时推送功能

- 在SecurityConfig中开放/ws/deviceStatus端点的匿名访问权限
- 新增DeviceStatusRedisListener实现Redis消息监听功能
- 新增DeviceStatusWebSocketHandler处理WebSocket连接和消息分发
- 新增RedisPubSubConfig配置Redis消息监听容器
- 新增SubscriptionManager管理WebSocket会话与设备的订阅关系
- 新增WebSocketConfig注册WebSocket处理器端点
- 实现设备状态变更的实时推送机制,支持批量订阅和取消订阅
- 添加心跳检测和错误处理机制确保连接稳定性
林仔 3 недель назад
Родитель
Сommit
335bcc9f04

+ 106 - 0
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/DeviceStatusRedisListener.java

@@ -0,0 +1,106 @@
+package com.zksy.web.websocket;
+
+import com.alibaba.fastjson.JSONObject;
+import com.zksy.base.domain.EquipmentStatus;
+import com.zksy.base.service.EquipmentStatusService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.data.redis.connection.Message;
+import org.springframework.data.redis.connection.MessageListener;
+import org.springframework.stereotype.Component;
+import org.springframework.web.socket.TextMessage;
+import org.springframework.web.socket.WebSocketSession;
+
+import java.io.IOException;
+import java.util.Set;
+
+/**
+ * Redis 设备状态变更监听器
+ * <p>
+ * 监听 Redis 频道 {@code device:status:change},当 zk-api-service 更新设备报警状态后,
+ * 接收消息 → 查询最新设备状态 → 推送到所有订阅了该设备的 WebSocket 会话。
+ *
+ * @author zksy
+ */
+@Slf4j
+@Component
+public class DeviceStatusRedisListener implements MessageListener {
+
+    @Autowired
+    private SubscriptionManager subscriptionManager;
+
+    @Autowired
+    private EquipmentStatusService equipmentStatusService;
+
+    /**
+     * Redis 频道名称
+     */
+    public static final String CHANNEL = "device:status:change";
+
+    @Override
+    public void onMessage(Message message, byte[] pattern) {
+        String body = new String(message.getBody());
+        log.debug("收到 Redis 设备状态变更消息: channel={}, body={}", CHANNEL, body);
+
+        try {
+            JSONObject msgObj = JSONObject.parseObject(body);
+            String deviceCode = msgObj.getString("deviceCode");
+
+            if (deviceCode == null || deviceCode.isEmpty()) {
+                log.warn("Redis 消息缺少 deviceCode 字段: {}", body);
+                return;
+            }
+
+            // 查找订阅了该设备的所有 WebSocket 会话
+            Set<String> sessionIds = subscriptionManager.getSubscribedSessions(deviceCode);
+            if (sessionIds.isEmpty()) {
+                log.debug("无会话订阅设备: deviceCode={}", deviceCode);
+                return;
+            }
+
+            // 查询最新设备状态
+            EquipmentStatus status = equipmentStatusService.getByDeviceCode(deviceCode);
+            if (status == null) {
+                log.warn("未找到设备状态: deviceCode={}", deviceCode);
+                return;
+            }
+
+            // 构建推送消息
+            JSONObject pushMsg = new JSONObject();
+            pushMsg.put("type", "deviceStatusPush");
+            pushMsg.put("data", status);
+
+            String pushPayload = pushMsg.toJSONString();
+            TextMessage textMessage = new TextMessage(pushPayload);
+
+            // 推送到所有订阅了该设备的会话
+            int successCount = 0;
+            int failCount = 0;
+            for (String sessionId : sessionIds) {
+                try {
+                    WebSocketSession session = DeviceStatusWebSocketHandler.getSession(sessionId);
+                    if (session != null && session.isOpen()) {
+                        synchronized (session) {
+                            session.sendMessage(textMessage);
+                        }
+                        successCount++;
+                    } else {
+                        // 会话已失效,清理订阅
+                        subscriptionManager.removeSession(sessionId);
+                        failCount++;
+                    }
+                } catch (IOException e) {
+                    log.error("推送设备状态失败: sessionId={}, deviceCode={}", sessionId, deviceCode, e);
+                    subscriptionManager.removeSession(sessionId);
+                    failCount++;
+                }
+            }
+
+            log.info("设备状态推送完成: deviceCode={}, 成功={}, 失败={}, 订阅会话数={}",
+                    deviceCode, successCount, failCount, sessionIds.size());
+
+        } catch (Exception e) {
+            log.error("处理 Redis 设备状态变更消息异常: body={}", body, e);
+        }
+    }
+}

+ 269 - 0
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/DeviceStatusWebSocketHandler.java

@@ -0,0 +1,269 @@
+package com.zksy.web.websocket;
+
+import com.alibaba.fastjson.JSONArray;
+import com.alibaba.fastjson.JSONObject;
+import com.zksy.base.domain.EquipmentStatus;
+import com.zksy.base.service.EquipmentStatusService;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+import org.springframework.web.socket.CloseStatus;
+import org.springframework.web.socket.TextMessage;
+import org.springframework.web.socket.WebSocketSession;
+import org.springframework.web.socket.handler.TextWebSocketHandler;
+
+import java.io.IOException;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+
+/**
+ * 设备状态 WebSocket 处理器(多路复用版)
+ * <p>
+ * 支持的消息类型(通过 JSON 中的 type 字段区分):
+ * <ul>
+ *   <li><b>subscribe</b> — 批量订阅设备,后端状态变更时主动推送</li>
+ *   <li><b>unsubscribe</b> — 取消订阅设备</li>
+ *   <li><b>deviceStatus</b> — 单次查询设备状态(兼容旧接口)</li>
+ *   <li><b>ping</b> — 心跳检测</li>
+ * </ul>
+ *
+ * <h3>协议示例</h3>
+ * <pre>
+ * // 订阅
+ * → {"type":"subscribe","deviceCodes":["DEV001","DEV002"]}
+ * ← {"type":"subscribed","deviceCodes":["DEV001","DEV002"]}
+ *
+ * // 取消订阅
+ * → {"type":"unsubscribe","deviceCodes":["DEV001"]}
+ * ← {"type":"unsubscribed","deviceCodes":["DEV001"]}
+ *
+ * // 心跳
+ * → {"type":"ping"}
+ * ← {"type":"pong","timestamp":1690000000000}
+ *
+ * // 单次查询(兼容旧接口)
+ * → {"type":"deviceStatus","deviceCode":"DEV001"}
+ * ← {"type":"deviceStatus","success":true,"data":{...}}
+ *
+ * // 设备状态变更推送(服务端主动推送)
+ * ← {"type":"deviceStatusPush","data":{...}}
+ * </pre>
+ *
+ * @author zksy
+ */
+@Slf4j
+@Component
+public class DeviceStatusWebSocketHandler extends TextWebSocketHandler {
+
+    @Autowired
+    private EquipmentStatusService equipmentStatusService;
+
+    @Autowired
+    private SubscriptionManager subscriptionManager;
+
+    /**
+     * 维护所有在线连接,key 为 sessionId
+     * 供 RedisListener 跨线程查找会话
+     */
+    private static final Map<String, WebSocketSession> SESSION_MAP = new ConcurrentHashMap<>();
+
+    /**
+     * 根据 sessionId 获取 WebSocket 会话(供 RedisListener 调用)
+     */
+    public static WebSocketSession getSession(String sessionId) {
+        return SESSION_MAP.get(sessionId);
+    }
+
+    @Override
+    public void afterConnectionEstablished(WebSocketSession session) {
+        String sessionId = session.getId();
+        SESSION_MAP.put(sessionId, session);
+        log.info("WebSocket 连接建立: sessionId={}, remoteAddress={}", sessionId, session.getRemoteAddress());
+    }
+
+    @Override
+    protected void handleTextMessage(WebSocketSession session, TextMessage message) throws IOException {
+        String payload = message.getPayload();
+        log.debug("收到 WebSocket 消息: sessionId={}, payload={}", session.getId(), payload);
+
+        try {
+            JSONObject request = JSONObject.parseObject(payload);
+            String type = request.getString("type");
+
+            if (type == null || type.isEmpty()) {
+                sendError(session, "消息缺少 type 字段");
+                return;
+            }
+
+            // === 根据 type 分发到不同的业务处理 ===
+            switch (type) {
+                //订阅
+                case "subscribe":
+                    handleSubscribe(session, request);
+                    break;
+                    //取消订阅
+                case "unsubscribe":
+                    handleUnsubscribe(session, request);
+                    break;
+                    //查询设备状态(单次查询)
+                case "deviceStatus":
+                    handleDeviceStatus(session, request);
+                    break;
+                    //心跳
+                case "ping":
+                    handlePing(session);
+                    break;
+                // 后续新增接口在此添加 case 分支
+                default:
+                    sendError(session, "不支持的消息类型: " + type);
+            }
+        } catch (Exception e) {
+            log.error("处理 WebSocket 消息异常: sessionId={}", session.getId(), e);
+            sendError(session, "消息处理失败: " + e.getMessage());
+        }
+    }
+
+    /**
+     * 处理批量订阅
+     */
+    private void handleSubscribe(WebSocketSession session, JSONObject request) throws IOException {
+        JSONArray deviceCodesArr = request.getJSONArray("deviceCodes");
+        if (deviceCodesArr == null || deviceCodesArr.isEmpty()) {
+            sendError(session, "缺少 deviceCodes 参数");
+            return;
+        }
+
+        Set<String> deviceCodes = new HashSet<>();
+        for (int i = 0; i < deviceCodesArr.size(); i++) {
+            String code = deviceCodesArr.getString(i);
+            if (code != null && !code.isEmpty()) {
+                deviceCodes.add(code);
+            }
+        }
+
+        if (deviceCodes.isEmpty()) {
+            sendError(session, "deviceCodes 不能为空");
+            return;
+        }
+
+        Set<String> subscribed = subscriptionManager.subscribe(session.getId(), deviceCodes);
+
+        JSONObject response = new JSONObject();
+        response.put("type", "subscribed");
+        response.put("deviceCodes", subscribed);
+        session.sendMessage(new TextMessage(response.toJSONString()));
+
+        log.info("订阅成功: sessionId={}, 订阅设备={}", session.getId(), subscribed);
+    }
+
+    /**
+     * 处理取消订阅
+     */
+    private void handleUnsubscribe(WebSocketSession session, JSONObject request) throws IOException {
+        JSONArray deviceCodesArr = request.getJSONArray("deviceCodes");
+        if (deviceCodesArr == null || deviceCodesArr.isEmpty()) {
+            sendError(session, "缺少 deviceCodes 参数");
+            return;
+        }
+
+        Set<String> deviceCodes = new HashSet<>();
+        for (int i = 0; i < deviceCodesArr.size(); i++) {
+            String code = deviceCodesArr.getString(i);
+            if (code != null && !code.isEmpty()) {
+                deviceCodes.add(code);
+            }
+        }
+
+        if (deviceCodes.isEmpty()) {
+            sendError(session, "deviceCodes 不能为空");
+            return;
+        }
+
+        Set<String> remaining = subscriptionManager.unsubscribe(session.getId(), deviceCodes);
+
+        JSONObject response = new JSONObject();
+        response.put("type", "unsubscribed");
+        response.put("deviceCodes", deviceCodes);
+        response.put("remaining", remaining);
+        session.sendMessage(new TextMessage(response.toJSONString()));
+
+        log.info("取消订阅: sessionId={}, 取消设备={}, 剩余订阅={}", session.getId(), deviceCodes, remaining);
+    }
+
+    /**
+     * 处理设备状态查询(兼容旧接口)
+     */
+    private void handleDeviceStatus(WebSocketSession session, JSONObject request) throws IOException {
+        String deviceCode = request.getString("deviceCode");
+        if (deviceCode == null || deviceCode.isEmpty()) {
+            sendResponse(session, "deviceStatus", false, null, "缺少 deviceCode 参数");
+            return;
+        }
+
+        log.info("查询设备状态: deviceCode={}", deviceCode);
+        EquipmentStatus status = equipmentStatusService.getByDeviceCode(deviceCode);
+
+        if (status == null) {
+            sendResponse(session, "deviceStatus", false, null, "未找到设备状态: deviceCode=" + deviceCode);
+        } else {
+            sendResponse(session, "deviceStatus", true, status, null);
+        }
+    }
+
+    /**
+     * 处理心跳
+     */
+    private void handlePing(WebSocketSession session) throws IOException {
+        JSONObject response = new JSONObject();
+        response.put("type", "pong");
+        response.put("timestamp", System.currentTimeMillis());
+        session.sendMessage(new TextMessage(response.toJSONString()));
+    }
+
+    @Override
+    public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
+        String sessionId = session.getId();
+        SESSION_MAP.remove(sessionId);
+        subscriptionManager.removeSession(sessionId);
+        log.info("WebSocket 连接关闭: sessionId={}, closeStatus={}", sessionId, status);
+    }
+
+    @Override
+    public void handleTransportError(WebSocketSession session, Throwable exception) {
+        String sessionId = session.getId();
+        SESSION_MAP.remove(sessionId);
+        subscriptionManager.removeSession(sessionId);
+        log.error("WebSocket 传输异常: sessionId={}", sessionId, exception);
+    }
+
+    /**
+     * 发送成功/失败响应
+     */
+    private void sendResponse(WebSocketSession session, String type, boolean success, Object data, String msg)
+            throws IOException {
+        JSONObject response = new JSONObject();
+        response.put("type", type);
+        response.put("success", success);
+        if (data != null) {
+            response.put("data", data);
+        }
+        if (msg != null) {
+            response.put("msg", msg);
+        }
+        session.sendMessage(new TextMessage(response.toJSONString()));
+    }
+
+    /**
+     * 发送错误响应
+     */
+    private void sendError(WebSocketSession session, String msg) throws IOException {
+        JSONObject response = new JSONObject();
+        response.put("type", "error");
+        response.put("success", false);
+        response.put("msg", msg);
+        session.sendMessage(new TextMessage(response.toJSONString()));
+    }
+}

+ 42 - 0
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/RedisPubSubConfig.java

@@ -0,0 +1,42 @@
+package com.zksy.web.websocket;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.data.redis.connection.RedisConnectionFactory;
+import org.springframework.data.redis.listener.ChannelTopic;
+import org.springframework.data.redis.listener.RedisMessageListenerContainer;
+
+/**
+ * Redis Pub/Sub 配置
+ * <p>
+ * 配置 Redis 消息监听容器,订阅设备状态变更频道。
+ * 当 zk-api-service 更新设备报警状态后,通过 Redis 频道通知本服务推送到 WebSocket 客户端。
+ *
+ * @author zksy
+ */
+@Configuration
+public class RedisPubSubConfig {
+
+    @Autowired
+    private DeviceStatusRedisListener deviceStatusRedisListener;
+
+    /**
+     * Redis 消息监听容器
+     * <p>
+     * 订阅 {@code device:status:change} 频道,收到消息后由
+     * {@link DeviceStatusRedisListener} 处理推送。
+     */
+    @Bean
+    public RedisMessageListenerContainer redisMessageListenerContainer(
+            RedisConnectionFactory connectionFactory) {
+        RedisMessageListenerContainer container = new RedisMessageListenerContainer();
+        container.setConnectionFactory(connectionFactory);
+
+        // 直接注册 MessageListener 实现,无需适配器
+        container.addMessageListener(deviceStatusRedisListener,
+                new ChannelTopic(DeviceStatusRedisListener.CHANNEL));
+
+        return container;
+    }
+}

+ 155 - 0
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/SubscriptionManager.java

@@ -0,0 +1,155 @@
+package com.zksy.web.websocket;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Component;
+
+import java.util.Collections;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+
+/**
+ * WebSocket 订阅管理器
+ * <p>
+ * 维护 WebSocket 会话与设备编码的双向映射关系,支持:
+ * <ul>
+ *   <li>sessionId → Set&lt;deviceCode&gt; 正向映射</li>
+ *   <li>deviceCode → Set&lt;sessionId&gt; 反向映射(快速查找需要推送的会话)</li>
+ * </ul>
+ * <p>
+ * 线程安全:使用 ConcurrentHashMap + ConcurrentHashMap.newKeySet() 保证并发安全。
+ *
+ * @author zksy
+ */
+@Slf4j
+@Component
+public class SubscriptionManager {
+
+    /**
+     * 正向映射:会话ID → 订阅的设备编码集合
+     */
+    private final ConcurrentHashMap<String, Set<String>> sessionToDevices = new ConcurrentHashMap<>();
+
+    /**
+     * 反向映射:设备编码 → 订阅了该设备的会话ID集合
+     */
+    private final ConcurrentHashMap<String, Set<String>> deviceToSessions = new ConcurrentHashMap<>();
+
+    /**
+     * 批量订阅设备
+     *
+     * @param sessionId   WebSocket 会话ID
+     * @param deviceCodes 要订阅的设备编码集合
+     * @return 当前会话订阅的所有设备编码
+     */
+    public Set<String> subscribe(String sessionId, Set<String> deviceCodes) {
+        if (deviceCodes == null || deviceCodes.isEmpty()) {
+            return getSubscribedDevices(sessionId);
+        }
+
+        // 更新正向映射
+        sessionToDevices.compute(sessionId, (k, existing) -> {
+            if (existing == null) {
+                existing = ConcurrentHashMap.newKeySet();
+            }
+            existing.addAll(deviceCodes);
+            return existing;
+        });
+
+        // 更新反向映射
+        for (String deviceCode : deviceCodes) {
+            deviceToSessions.compute(deviceCode, (k, existing) -> {
+                if (existing == null) {
+                    existing = ConcurrentHashMap.newKeySet();
+                }
+                existing.add(sessionId);
+                return existing;
+            });
+        }
+
+        log.info("订阅成功: sessionId={}, deviceCodes={}, 当前订阅总数={}",
+                sessionId, deviceCodes, getSubscribedDevices(sessionId).size());
+        return getSubscribedDevices(sessionId);
+    }
+
+    /**
+     * 取消订阅设备
+     *
+     * @param sessionId   WebSocket 会话ID
+     * @param deviceCodes 要取消订阅的设备编码集合
+     * @return 当前会话剩余的订阅设备编码
+     */
+    public Set<String> unsubscribe(String sessionId, Set<String> deviceCodes) {
+        if (deviceCodes == null || deviceCodes.isEmpty()) {
+            return getSubscribedDevices(sessionId);
+        }
+
+        // 更新正向映射
+        sessionToDevices.computeIfPresent(sessionId, (k, existing) -> {
+            existing.removeAll(deviceCodes);
+            return existing.isEmpty() ? null : existing;
+        });
+
+        // 更新反向映射
+        for (String deviceCode : deviceCodes) {
+            deviceToSessions.computeIfPresent(deviceCode, (k, existing) -> {
+                existing.remove(sessionId);
+                return existing.isEmpty() ? null : existing;
+            });
+        }
+
+        log.info("取消订阅: sessionId={}, deviceCodes={}, 当前订阅总数={}",
+                sessionId, deviceCodes, getSubscribedDevices(sessionId).size());
+        return getSubscribedDevices(sessionId);
+    }
+
+    /**
+     * 移除会话的所有订阅(连接断开时调用)
+     *
+     * @param sessionId WebSocket 会话ID
+     */
+    public void removeSession(String sessionId) {
+        Set<String> deviceCodes = sessionToDevices.remove(sessionId);
+        if (deviceCodes == null || deviceCodes.isEmpty()) {
+            return;
+        }
+
+        for (String deviceCode : deviceCodes) {
+            deviceToSessions.computeIfPresent(deviceCode, (k, existing) -> {
+                existing.remove(sessionId);
+                return existing.isEmpty() ? null : existing;
+            });
+        }
+
+        log.info("会话移除,清理订阅: sessionId={}, 清理设备数={}", sessionId, deviceCodes.size());
+    }
+
+    /**
+     * 获取订阅了指定设备的所有会话ID
+     *
+     * @param deviceCode 设备编码
+     * @return 会话ID集合(不可变),未找到返回空集合
+     */
+    public Set<String> getSubscribedSessions(String deviceCode) {
+        Set<String> sessions = deviceToSessions.get(deviceCode);
+        return sessions != null ? Collections.unmodifiableSet(sessions) : Collections.emptySet();
+    }
+
+    /**
+     * 获取指定会话订阅的所有设备编码
+     *
+     * @param sessionId WebSocket 会话ID
+     * @return 设备编码集合(不可变),未找到返回空集合
+     */
+    public Set<String> getSubscribedDevices(String sessionId) {
+        Set<String> devices = sessionToDevices.get(sessionId);
+        return devices != null ? Collections.unmodifiableSet(devices) : Collections.emptySet();
+    }
+
+    /**
+     * 获取当前订阅统计信息
+     */
+    public String getStats() {
+        return String.format("活跃会话数=%d, 被订阅设备数=%d",
+                sessionToDevices.size(), deviceToSessions.size());
+    }
+}

+ 30 - 0
pipe-network-service/zksy-admin/src/main/java/com/zksy/web/websocket/WebSocketConfig.java

@@ -0,0 +1,30 @@
+package com.zksy.web.websocket;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.web.socket.config.annotation.EnableWebSocket;
+import org.springframework.web.socket.config.annotation.WebSocketConfigurer;
+import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry;
+
+/**
+ * WebSocket 配置类
+ * <p>
+ * 注册 WebSocket 处理器,供可视化等前端调用。
+ * 后续新增接口只需在此注册新的 handler 即可。
+ *
+ * @author zksy
+ */
+@Configuration
+@EnableWebSocket
+public class WebSocketConfig implements WebSocketConfigurer {
+
+    @Autowired
+    private DeviceStatusWebSocketHandler deviceStatusWebSocketHandler;
+
+    @Override
+    public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
+        // 设备状态查询 WebSocket 端点
+        registry.addHandler(deviceStatusWebSocketHandler, "/ws/deviceStatus")
+                .setAllowedOrigins("*");
+    }
+}

+ 1 - 1
pipe-network-service/zksy-framework/src/main/java/com/zksy/framework/config/SecurityConfig.java

@@ -111,7 +111,7 @@ public class SecurityConfig
             .authorizeHttpRequests((requests) -> {
                 permitAllUrl.getUrls().forEach(url -> requests.antMatchers(url).permitAll());
                 // 对于登录login 注册register 验证码captchaImage 允许匿名访问
-                requests.antMatchers("/login", "/register", "/captchaImage","/**/visualization/**").permitAll()
+                requests.antMatchers("/login", "/register", "/captchaImage","/**/visualization/**","/ws/deviceStatus").permitAll()
                     // 静态资源,可匿名访问
                     .antMatchers(HttpMethod.GET, "/", "/*.html", "/**/*.html", "/**/*.css", "/**/*.js", "/profile/**").permitAll()
                     .antMatchers("/swagger-ui.html", "/swagger-resources/**", "/webjars/**", "/*/api-docs", "/druid/**").permitAll()