From f4f4fc53260eb67483dce406a963628273786a61 Mon Sep 17 00:00:00 2001
From: liusuyi <1951119284@qq.com>
Date: Sat, 30 May 2026 17:22:17 +0800
Subject: [PATCH] 优化

---
 ard-modules/ard-modules-work/src/main/java/com/ard/work/websocket/utils/PTZWebSocketUtils.java |   93 +++++++++++++++++-----------------------------
 1 files changed, 35 insertions(+), 58 deletions(-)

diff --git a/ard-modules/ard-modules-work/src/main/java/com/ard/work/websocket/utils/PTZWebSocketUtils.java b/ard-modules/ard-modules-work/src/main/java/com/ard/work/websocket/utils/PTZWebSocketUtils.java
index 905f768..3ae10a5 100644
--- a/ard-modules/ard-modules-work/src/main/java/com/ard/work/websocket/utils/PTZWebSocketUtils.java
+++ b/ard-modules/ard-modules-work/src/main/java/com/ard/work/websocket/utils/PTZWebSocketUtils.java
@@ -12,6 +12,7 @@
 import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.Executor;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
 
@@ -25,13 +26,15 @@
 @Slf4j
 public final class PTZWebSocketUtils {
 
-    // 存储 websocket session
     public static final ConcurrentMap<String, WebSocketSession> ONLINE_USER_SESSIONS = new ConcurrentHashMap<>();
 
-    /**
-     * @param session 用户 session
-     * @param message 发送内容
-     */
+    /** 推送线程池,由 PTZWebSocketInitializer 在启动时注入;未注入时降级为同步执行 */
+    private static volatile Executor pushExecutor;
+
+    public static void setPushExecutor(Executor executor) {
+        pushExecutor = executor;
+    }
+
     public static void sendMessage(WebSocketSession session, String message) {
         if (session == null) {
             return;
@@ -42,84 +45,58 @@
         }
         synchronized (session) {
             try {
-                //log.debug("发送消息:" + message);
-                session.sendMessage(new TextMessage(String.join(", ", message)));
+                session.sendMessage(new TextMessage(message));
             } catch (IOException e) {
                 log.error("sendMessage IOException ", e);
             }
         }
     }
 
-    /**
-     * @param session 用户 session
-     * @param message 发送内容
-     */
     public static void sendMessage(WebSocketSession session, Map message) {
-        if (session == null) {
-            return;
-        }
-       final InetSocketAddress remoteAddress = session.getRemoteAddress();
-        if (remoteAddress == null) {
-            return;
-        }
-        synchronized (session) {
-            try {
-                session.sendMessage(new TextMessage(JSONObject.toJSONString(message)));
-            } catch (IOException e) {
-                log.error("sendMessage IOException ", e);
-            }
-        }
+        sendMessage(session, JSONObject.toJSONString(message));
     }
 
     public static void sendMessage(WebSocketSession session, List message) {
-        if (session == null) {
-            return;
-        }
-        final InetSocketAddress remoteAddress = session.getRemoteAddress();
-        if (remoteAddress == null) {
-            return;
-        }
-        synchronized (session) {
-            try {
-                session.sendMessage(new TextMessage(JSONObject.toJSONString(message)));
-            } catch (IOException e) {
-                log.error("sendMessage IOException ", e);
-            }
-        }
+        sendMessage(session, JSON.toJSONString(message));
     }
 
     /**
-     * 推送消息到其他客户端
-     *
-     * @param message
+     * 推送字符串消息到所有客户端
      */
     public static void sendMessageAll(String message) {
-        ONLINE_USER_SESSIONS.forEach((sessionId, session) -> sendMessage(session, message));
+        broadcast(message);
     }
 
     /**
-     * 推送消息到其他客户端
-     *
-     * @param message
+     * 推送Map消息到所有客户端
      */
     public static void sendMessageAll(Map message) {
-        JSONObject jsonObject = new JSONObject(message);
-        ONLINE_USER_SESSIONS.forEach((sessionId, session) -> sendMessage(session, jsonObject.toString()));
+        broadcast(JSONObject.toJSONString(message));
     }
+
     /**
-     * 推送消息到其他客户端
-     *
-     * @param message
+     * 推送List消息到所有客户端
      */
     public static void sendMessageAll(List message) {
-        ONLINE_USER_SESSIONS.forEach((sessionId, session) -> sendMessage(session, JSON.toJSONString(message)));
+        broadcast(JSON.toJSONString(message));
     }
+
     /**
-     * @Author 刘苏义
-     * @Description  发送消息给userId_前缀的人
-     * @Date   2024/7/16 10:24
-     * @Param
-     * @return
+     * 一次序列化,并行分发到所有 session
+     */
+    private static void broadcast(String payload) {
+        if (ONLINE_USER_SESSIONS.isEmpty()) return;
+        Executor executor = pushExecutor;
+        if (executor == null) {
+            ONLINE_USER_SESSIONS.values().forEach(session -> sendMessage(session, payload));
+            return;
+        }
+        ONLINE_USER_SESSIONS.values().forEach(session ->
+                executor.execute(() -> sendMessage(session, payload)));
+    }
+
+    /**
+     * 发送消息给 userId_ 前缀的人
      */
     public static void sendMessagePrefix(String targetId, String message) {
         String regex = "^" + Pattern.quote(targetId) + "_\\d+$";
@@ -129,4 +106,4 @@
             return matcher.matches();
         }).map(Map.Entry::getValue).forEach(session -> PTZWebSocketUtils.sendMessage(session, message));
     }
-}
\ No newline at end of file
+}

--
Gitblit v1.9.3