websocket消息服务器模块

This commit is contained in:
tangw 2025-11-27 19:07:55 +08:00
parent 289187847f
commit dbea5a7b3f
9 changed files with 829 additions and 0 deletions

View File

@ -0,0 +1,295 @@
package com.czlis.interfaceCommon.websocket;
import com.alibaba.fastjson2.JSON;
import com.czlis.common.core.domain.model.LoginUser;
import com.czlis.common.utils.SecurityUtils;
import com.czlis.common.utils.spring.SpringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.socket.server.standard.SpringConfigurator;
import javax.websocket.*;
import javax.websocket.server.PathParam;
import javax.websocket.server.ServerEndpoint;
import java.net.InetSocketAddress;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
/**
* lis消息中心服务器
* 原理:核心是 统一管理所有在线 WebSocket 连接(而非直接引用实例)—— 因为 @ServerEndpoint 是「多实例」模式(每个客户端连接对应一个独立实例),直接引用单个实例毫无意义,且会导致重复创建 / 连接泄漏问题。
* 业务结构:在 WebSocket 端点类内部维护一个「全局共享的在线连接容器」,其他业务通过这个容器获取目标连接并发送消息
* 一、核心设计原则
* 全局连接容器:用线程安全的集合(如 ConcurrentHashMap)存储所有在线 Session(WebSocket 连接会话,每个客户端唯一),避免多线程并发问题;
* 实例无状态:@ServerEndpoint 实例仅用于处理连接生命周期,不存储业务数据,所有共享状态(如在线连接)放在「静态容器」中;
* 工具方法暴露:在 WebSocket 端点类中提供静态工具方法(如 sendToUser、broadcast),其他业务直接调用该方法,无需关心实例创建;
* 避免重复连接:连接的创建 / 销毁由 WebSocket 协议自动管理(@OnOpen 新增、@OnClose 移除),业务层不干预实例创建。
*/
@ServerEndpoint(value = "/ws/msgserver/{userId}", configurator = SpringConfigurator.class)
public class ChatWebSocketServer {
// 🌟 核心固定参数白名单:这些参数不存入动态属性(单独解析为固定身份属性)
private static final Set<String> CORE_PARAMS = new HashSet<>(Arrays.asList("token", "username", "role"));
// 🌟 容器改为存储封装后的WebSocketUser(key=userId)
private static final Map<String, WebSocketUser> ONLINE_USER_MAP = new ConcurrentHashMap<>();
/**
* 连接建立时:封装用户信息并存入容器
* 注意:若需要传递username/role等属性,可通过路径参数/请求参数/Token解析
*/
@OnOpen
public void onOpen(@PathParam("userId") String userId, Session session) {
// 1. 获取所有请求参数(客户端连接时携带的参数,如 ws://xxx/ws/chat/1001?token=xxx&username=张三&role=USER&deviceType=pc&version=1.0)
Map<String, List<String>> requestParams = session.getRequestParameterMap();
String username = "";
String role = "";
String loginIp = getClientIp(session);
String token="";
try {
// 场景1:正常业务场景(前端登录后连接)→ 从 SecurityUtils 获取
LoginUser loginUser = SecurityUtils.getLoginUser();
token=loginUser.getToken();
username = loginUser.getUser().getNickName();
role = loginUser.getUser().getRoleIds().toString();
} catch (Exception e) {
// 场景2:测试工具场景(无登录上下文)→ 从 URL 参数解析兜底
// 从参数中获取 username/role(测试工具连接时手动拼接)
username = requestParams.getOrDefault("username", Collections.singletonList("测试用户")).get(0);
role = requestParams.getOrDefault("role", Collections.singletonList("TEST")).get(0);
System.out.println("未获取到登录上下文,使用URL参数兜底:用户名=" + username + ",角色=" + role);
}
// 3. 封装核心用户信息(固定属性)
WebSocketUser webSocketUser = new WebSocketUser(userId, username, role, loginIp, session);
webSocketUser.setDynamicAttr("token", token);
// 4. 🌟 自动解析剩余参数,存入dynamicAttrs(过滤核心参数)
for (Map.Entry<String, List<String>> entry : requestParams.entrySet()) {
String paramKey = entry.getKey();
// 跳过核心固定参数、空值参数
if (CORE_PARAMS.contains(paramKey) || entry.getValue() == null || entry.getValue().isEmpty()) {
continue;
}
// 参数值取第一个(WebSocket请求参数通常是单值,多值可自行处理)
String paramValue = entry.getValue().get(0);
// 4. 跳过空字符串/无效值(可选,进一步过滤垃圾数据)
if (paramValue.isEmpty() || "null".equalsIgnoreCase(paramValue)) {
continue;
}
// 自动存入动态属性
webSocketUser.setDynamicAttr(paramKey, paramValue);
}
// 5. 存入全局容器
ONLINE_USER_MAP.put(userId, webSocketUser);
System.out.println("用户[" + userId + "-" + username + "]连接成功,IP:" + loginIp + ",当前在线人数:" + ONLINE_USER_MAP.size());
}
// ========== 消息接收核心方法 ==========
/**
* 业务层处理客户端通过 WebSocket 发送的消息,核心思路是:将「消息接收」与「业务处理」解耦 —— WebSocket 端点类仅负责接收消息、做基础校验,再将消息转发给业务层处理(同步 / 异步),业务层专注于业务逻辑(如消息解析、权限校验、数据处理、结果推送),最终通过 WebSocket 端点类将处理结果返回给客户端。
* 整体流程设计
* 客户端发送消息 → WebSocket端点(@OnMessage)接收 → 基础校验 → 转发业务层处理 → 业务层执行业务逻辑 → 结果通过WebSocket推送给客户端
*
* @param message
* @param session
* @param userId
*/
@OnMessage
public void onMessage(String message, Session session, @PathParam("userId") String userId) {
try {
// 1. 基础校验:消息为空/连接关闭 → 直接返回
if (message == null || message.trim().isEmpty() || !session.isOpen()) {
return;
}
// 2. 解析消息为统一协议对象(JSON解析)
WebSocketMessage wsMessage = JSON.parseObject(message, WebSocketMessage.class);
if (wsMessage == null) {
sendToUser(userId, "错误:消息格式不合法,请按JSON格式发送");
return;
}
//暂时使用默认时间戳,生产环境应该使用客户端必传字段
if (wsMessage.getTimestamp() == null || wsMessage.getTimestamp() == 0) {
wsMessage.setTimestamp(System.currentTimeMillis());
}
// 3. 补充发送者ID(从路径参数获取,避免客户端伪造)
wsMessage.setSenderId(userId);
// 4. 获取Spring管理的业务层Bean,转发消息处理
WebSocketMsgService msgService = SpringUtils.getBean(WebSocketMsgService.class);
if (msgService == null) {
throw new RuntimeException("获取WebSocketMsgService Bean失败,请检查Spring配置");
}
// 同步处理(简单场景)/ 异步处理(高并发场景,推荐)
msgService.handleMessage(wsMessage, session);
} catch (Exception e) {
e.printStackTrace();
// 异常反馈给客户端
RetryMsgService retryMsgService = SpringUtils.getBean(RetryMsgService.class);
// 封装错误消息存入重试队列
WebSocketMessage errorMsg = new WebSocketMessage();
errorMsg.setSenderId(userId.toString());
errorMsg.setMsgType("ERROR");
errorMsg.setData("消息处理失败:" + e.getMessage());
errorMsg.setTimestamp(System.currentTimeMillis());
retryMsgService.sendToRetryQueue(errorMsg);
}
}
/**
* 连接关闭时:移除用户信息
*/
@OnClose
public void onClose(@PathParam("userId") String userId, Session session) {
ONLINE_USER_MAP.remove(userId);
System.out.println("用户[" + userId + "]断开连接,当前在线人数:" + ONLINE_USER_MAP.size());
}
/**
* 通信异常时:移除用户信息
*/
@OnError
public void onError(@PathParam("userId") String userId, Session session, Throwable throwable) {
throwable.printStackTrace();
ONLINE_USER_MAP.remove(userId);
}
// ========== 工具方法:获取客户端IP ==========
private String getClientIp(Session session) {
try {
// 获取客户端远程地址(不同服务器容器可能有差异)
InetSocketAddress remoteAddress = (InetSocketAddress) session.getUserProperties().get("javax.websocket.endpoint.remoteAddress");
return remoteAddress.getAddress().getHostAddress();
} catch (Exception e) {
return "未知IP";
}
}
// ========== 🌟 暴露属性查询方法(业务层调用) ==========
/**
* 获取所有在线用户的完整信息
*/
public static Collection<WebSocketUser> getAllOnlineUsers() {
// 返回快照,避免并发修改
return Collections.unmodifiableCollection(new ArrayList<>(ONLINE_USER_MAP.values()));
}
/**
* 根据userId获取用户完整信息
*/
public static WebSocketUser getOnlineUser(String userId) {
return ONLINE_USER_MAP.get(userId);
}
/**
* 根据角色筛选在线用户(如仅获取管理员)
*/
public static List<WebSocketUser> getOnlineUsersByRole(String role) {
List<WebSocketUser> roleUsers = new ArrayList<>();
for (WebSocketUser user : ONLINE_USER_MAP.values()) {
if (role.equals(user.getRole()) && user.isConnected()) {
roleUsers.add(user);
}
}
return roleUsers;
}
/**
* 获取在线用户的userId列表(兼容原有逻辑)
*/
public static Set<String> getAllOnlineUserIds() {
return Collections.unmodifiableSet(new HashSet<>(ONLINE_USER_MAP.keySet()));
}
// ========== 消息发送方法(适配新容器) ==========
public static boolean sendToUser(String userId, String message) {
WebSocketUser user = ONLINE_USER_MAP.get(userId);
if (user == null || !user.isConnected()) {
System.out.println("用户[" + userId + "]不在线或连接已关闭");
return false;
}
try {
user.getSession().getAsyncRemote().sendText(message);
return true;
} catch (Exception e) {
e.printStackTrace();
ONLINE_USER_MAP.remove(userId);
return false;
}
}
public static void broadcast(String message) {
for (WebSocketUser user : ONLINE_USER_MAP.values()) {
if (user.isConnected()) {
user.getSession().getAsyncRemote().sendText(message);
}
}
}
// ========== 🌟 封装属性修改方法(仅允许修改动态属性) ==========
/**
* 给指定用户设置单个动态属性(业务层调用)
*
* @param userId 用户ID
* @param attrKey 属性键
* @param attrValue 属性值
* @return true=修改成功(用户在线),false=用户不在线
*/
public static boolean setUserDynamicAttr(String userId, String attrKey, Object attrValue) {
// 1. 校验用户是否在线
WebSocketUser user = ONLINE_USER_MAP.get(userId);
if (user == null || !user.isConnected()) {
System.out.println("用户[" + userId + "]不在线,无法修改属性");
return false;
}
// 2. 禁止修改核心属性(防止业务层误操作)
if (CORE_PARAMS.contains(attrKey)) {
System.out.println("禁止修改核心属性:" + attrKey);
return false;
}
// 3. 受控修改动态属性(WebSocketUser的dynamicAttrs是线程安全的)
user.setDynamicAttr(attrKey, attrValue);
System.out.println("用户[" + userId + "]的属性[" + attrKey + "]已更新为:" + attrValue);
return true;
}
/**
* 批量修改指定用户的动态属性
*
* @param userId 用户ID
* @param attrMap 待修改的属性键值对
* @return true=修改成功
*/
public static boolean batchSetUserDynamicAttr(String userId, Map<String, Object> attrMap) {
WebSocketUser user = ONLINE_USER_MAP.get(userId);
if (user == null || !user.isConnected()) {
return false;
}
// 批量设置(过滤核心属性)
for (Map.Entry<String, Object> entry : attrMap.entrySet()) {
String attrKey = entry.getKey();
if (!CORE_PARAMS.contains(attrKey)) {
user.setDynamicAttr(attrKey, entry.getValue());
}
}
return true;
}
/**
* 移除指定用户的某个动态属性
*
* @param userId 用户ID
* @param attrKey 属性键
* @return true=移除成功
*/
public static boolean removeUserDynamicAttr(String userId, String attrKey) {
WebSocketUser user = ONLINE_USER_MAP.get(userId);
if (user == null || !user.isConnected()) {
return false;
}
user.removeDynamicAttr(attrKey);
return true;
}
}

View File

@ -0,0 +1,107 @@
package com.czlis.interfaceCommon.websocket;
import com.alibaba.fastjson2.JSON;
import com.czlis.common.core.redis.RedisCache;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.util.concurrent.TimeUnit;
/**
* 消息重试服务(若依框架适配)
*/
@Service
public class RetryMsgService {
// Redis键前缀(若依框架建议添加业务前缀,避免键冲突)
private static final String RETRY_QUEUE_KEY = "ws:msg:retry:queue";
private static final String DEAD_LETTER_KEY = "ws:msg:dead:letter";
private static final String RETRY_COUNT_KEY_PREFIX = "ws:msg:retry:count:";
// 最大重试次数(可配置在若依的application.yml中,用@Value注入)
private static final int MAX_RETRY_COUNT = 3;
@Autowired
private RedisCache redisCache;
@Autowired
private WebSocketMsgService webSocketMsgService;
/**
* 存入重试队列(使用若依RedisTemplate)
*/
public void sendToRetryQueue(WebSocketMessage wsMessage) {
if (wsMessage == null) {
return;
}
String msgId = generateMsgId(wsMessage);
String msgJson = JSON.toJSONString(wsMessage);
// 若依Redis:存储重试次数
redisCache.setCacheObject(RETRY_COUNT_KEY_PREFIX + msgId, "0", 24, TimeUnit.HOURS);
// 存入重试队列
redisCache.redisTemplate.opsForList().leftPush(RETRY_QUEUE_KEY, msgJson);
System.out.println("消息[" + msgId + "]存入重试队列");
}
/**
* 定时重试处理(由定时任务调用)
*/
public void retryProcessMsg() {
// 从队列取出消息(阻塞1秒)
String msgJson = (String) redisCache.redisTemplate.opsForList().rightPop(
RETRY_QUEUE_KEY,
1,
TimeUnit.SECONDS
);
if (msgJson == null) {
return;
}
try {
WebSocketMessage wsMessage = JSON.parseObject(msgJson, WebSocketMessage.class);
String msgId = generateMsgId(wsMessage);
// 🌟 替换方案2:使用RedisCache封装的getCacheObject方法(更简洁)
String retryCountStr = redisCache.getCacheObject(RETRY_COUNT_KEY_PREFIX + msgId);
int retryCount = retryCountStr == null ? 0 : Integer.parseInt(retryCountStr);
// 超过最大重试次数,转入死信队列
if (retryCount >= MAX_RETRY_COUNT) {
redisCache.redisTemplate.opsForList().leftPush(DEAD_LETTER_KEY, msgJson);
redisCache.redisTemplate.delete(RETRY_COUNT_KEY_PREFIX + msgId);
System.out.println("消息[" + msgId + "]重试" + MAX_RETRY_COUNT + "次失败,转入死信队列");
return;
}
// 重试处理
webSocketMsgService.handleMessage(wsMessage, null);
// 处理成功,删除重试计数
redisCache.redisTemplate.delete(RETRY_COUNT_KEY_PREFIX + msgId);
System.out.println("消息[" + msgId + "]重试第" + (retryCount + 1) + "次成功");
} catch (Exception e) {
e.printStackTrace();
WebSocketMessage wsMessage = JSON.parseObject(msgJson, WebSocketMessage.class);
String msgId = generateMsgId(wsMessage);
int retryCount = Integer.parseInt(redisCache.getCacheObject(RETRY_COUNT_KEY_PREFIX + msgId));
retryCount++;
// 更新重试次数,重新存入队列
redisCache.setCacheObject(
RETRY_COUNT_KEY_PREFIX + msgId,
String.valueOf(retryCount),
24,
TimeUnit.HOURS
);
redisCache.redisTemplate.opsForList().leftPush(RETRY_QUEUE_KEY, msgJson);
System.out.println("消息[" + msgId + "]重试第" + retryCount + "次失败,重新存入队列");
}
}
/**
* 生成消息唯一标识
*/
private String generateMsgId(WebSocketMessage wsMessage) {
return wsMessage.getSenderId() + "_" + wsMessage.getTimestamp();
}
}

View File

@ -0,0 +1,21 @@
package com.czlis.interfaceCommon.websocket;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
@Component
public class RetryMsgTask {
@Autowired
private RetryMsgService retryMsgService;
/**
* 每5秒执行一次(若依框架支持cron表达式或fixedRate)
*/
@Scheduled(fixedRate = 5000)
public void executeRetry() {
// 多线程处理,避免阻塞定时任务线程
new Thread(retryMsgService::retryProcessMsg).start();
}
}

View File

@ -0,0 +1,60 @@
package com.czlis.interfaceCommon.websocket;
import org.springframework.stereotype.Service;
import java.util.List;
@Service
public class UserService {
/**
* 示例1:获取所有在线管理员并推送通知
*/
public void pushNoticeToAdmin(String notice) {
// 根据角色筛选在线管理员
List<WebSocketUser> adminUsers = ChatWebSocketServer.getOnlineUsersByRole("ADMIN");
if (adminUsers.isEmpty()) {
System.out.println("当前无在线管理员");
return;
}
String message = "【管理员通知】" + notice;
for (WebSocketUser admin : adminUsers) {
ChatWebSocketServer.sendToUser(admin.getUserId(), message);
System.out.println("已给管理员[" + admin.getUsername() + "]推送通知,IP:" + admin.getLoginIp());
}
}
/**
* 示例2:获取指定用户的详细连接信息
*/
public void getUserConnectInfo(String userId) {
WebSocketUser user = ChatWebSocketServer.getOnlineUser(userId);
if (user == null) {
System.out.println("用户[" + userId + "]不在线");
return;
}
// 获取身份属性
System.out.println("用户名:" + user.getUsername());
System.out.println("角色:" + user.getRole());
System.out.println("登录IP:" + user.getLoginIp());
System.out.println("连接时间:" + user.getConnectTime());
System.out.println("是否在线:" + user.isConnected());
// 获取动态属性(如设备类型)
String deviceType = (String) user.getDynamicAttr("deviceType");
System.out.println("登录设备:" + (deviceType == null ? "未知" : deviceType));
}
/**
* 示例3:给在线用户设置动态属性(如设备类型)
*/
public void setUserDeviceAttr(String userId, String deviceType) {
WebSocketUser user = ChatWebSocketServer.getOnlineUser(userId);
if (user != null) {
user.setDynamicAttr("deviceType", deviceType);
System.out.println("已给用户[" + userId + "]设置设备类型:" + deviceType);
}
}
}

View File

@ -0,0 +1,24 @@
package com.czlis.interfaceCommon.websocket;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
* 客户端与服务端的统一消息协议
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
public class WebSocketMessage {
// 消息类型(如CHAT/NOTICE/ORDER/HEARTBEAT)
private String msgType;
// 业务数据(JSON字符串,可解析为对应业务DTO)
private String data;
// 发送者ID(可选,也可从WebSocket连接中获取)
private String senderId;
// 接收者ID(可选,单聊/广播区分)
private String receiverId;
// 消息时间戳
private Long timestamp;
}

View File

@ -0,0 +1,15 @@
package com.czlis.interfaceCommon.websocket;
import javax.websocket.Session;
/**
* WebSocket消息处理业务接口
*/
public interface WebSocketMsgService {
/**
* 处理客户端发送的消息
* @param wsMessage 解析后的统一消息对象
* @param session 客户端WebSocket会话
*/
void handleMessage(WebSocketMessage wsMessage, Session session);
}

View File

@ -0,0 +1,172 @@
package com.czlis.interfaceCommon.websocket;
import com.alibaba.fastjson2.JSON;
import com.czlis.common.annotation.Async;
import com.czlis.common.core.redis.RedisCache;
import com.czlis.common.utils.spring.SpringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Service;
import javax.websocket.Session;
import java.util.Map;
import java.util.concurrent.Executor;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
/**业务层(Spring Bean)专注于「消息解析、权限校验、业务逻辑处理、结果推送」,支持不同类型消息的差异化处理。
* WebSocket消息处理业务实现类
*/
@Service
public class WebSocketMsgServiceImpl implements WebSocketMsgService {
@Override
public void handleMessage(WebSocketMessage wsMessage, Session session) {
if (wsMessage.getTimestamp() == null || wsMessage.getTimestamp() == 0) {
ChatWebSocketServer.sendToUser(wsMessage.getSenderId(), "错误:消息必须携带时间戳");
return;
}
String msgId = wsMessage.getSenderId() + "_" + wsMessage.getTimestamp();
// 校验消息是否已处理(如Redis中是否存在msgId)
if (SpringUtils.getBean(RedisCache.class).redisTemplate.hasKey("ws:msg:processed:" + msgId)) {
ChatWebSocketServer.sendToUser(wsMessage.getSenderId(), "提示:消息已处理,无需重复发送");
return;
}
try {
// 1. 基础校验:消息类型、发送者ID
String msgType = wsMessage.getMsgType();
String senderId = wsMessage.getSenderId();
if (msgType == null || senderId == null) {
ChatWebSocketServer.sendToUser(senderId, "错误:消息类型/发送者ID不能为空");
return;
}
// 2. 根据消息类型分发处理(策略模式,可扩展)
switch (msgType) {
case "CHAT": // 单聊消息
handleChatMsg(wsMessage);
break;
case "BROADCAST": // 广播消息
handleBroadcastMsg(wsMessage);
break;
case "ORDER": // 订单相关消息
handleOrderMsg(wsMessage);
break;
case "HEARTBEAT": // 心跳包(仅确认,无需处理业务)
handleHeartbeatMsg(wsMessage);
break;
default:
ChatWebSocketServer.sendToUser(senderId, "错误:不支持的消息类型:" + msgType);
}
} catch (Exception e) {
e.printStackTrace();
// ChatWebSocketServer.sendToUser(wsMessage.getSenderId(), "错误:消息处理失败:" + e.getMessage());
String senderId = wsMessage.getSenderId();
if (senderId != null && session != null && session.isOpen()) {
ChatWebSocketServer.sendToUser(senderId, "错误:消息处理失败,已加入重试队列");
}
RetryMsgService retryMsgService = SpringUtils.getBean(RetryMsgService.class);
retryMsgService.sendToRetryQueue(wsMessage);
}
// 标记消息已处理(设置过期时间,避免内存泄漏)
SpringUtils.getBean(RedisCache.class).redisTemplate.opsForValue().set("ws:msg:processed:" + msgId, "1", 24, TimeUnit.HOURS);
}
/**
* 处理单聊消息
*/
private void handleChatMsg(WebSocketMessage wsMessage) {
/*
String senderId = wsMessage.getSenderId();
String receiverId = wsMessage.getReceiverId();
String data = wsMessage.getData();
// 1. 校验接收者是否在线
boolean receiverOnline = ChatWebSocketServer.getOnlineUser(receiverId) != null;
if (!receiverOnline) {
ChatWebSocketServer.sendToUser(senderId, "提示:接收者[" + receiverId + "]不在线,消息已存入离线队列");
// 调用离线消息服务:存入Redis/MQ,等待接收者上线后推送
// userService.saveOfflineMessage(receiverId, data);
return;
}
// 2. 组装转发消息(添加发送者信息)
Map<String, Object> chatData = JSON.parseObject(data);
// chatData.put("senderName", userService.getUsernameById(senderId)); // 获取发送者昵称
String forwardMsg = JSON.toJSONString(chatData);
// 3. 转发消息给接收者
boolean sendSuccess = ChatWebSocketServer.sendToUser(receiverId, forwardMsg);
if (sendSuccess) {
// 给发送者反馈:消息已送达
ChatWebSocketServer.sendToUser(senderId, "提示:消息已送达接收者[" + receiverId + "]");
}
*/
}
/**
* 处理广播消息(仅管理员可发送)
*/
private void handleBroadcastMsg(WebSocketMessage wsMessage) {
/*
String senderId = wsMessage.getSenderId();
String data = wsMessage.getData();
// 1. 权限校验:仅管理员可发送广播
boolean isAdmin = "ADMIN".equals(userService.getUserRoleById(senderId));
if (!isAdmin) {
ChatWebSocketServer.sendToUser(senderId, "错误:无广播消息权限");
return;
}
// 2. 广播消息给所有在线用户
String broadcastMsg = "【系统广播】" + data;
ChatWebSocketServer.broadcast(broadcastMsg);
// 3. 给发送者反馈
ChatWebSocketServer.sendToUser(senderId, "提示:广播消息已发送,覆盖在线人数:" + ChatWebSocketServer.getOnlineUserCount());
*/
}
/**
* 处理订单相关消息(如查询订单、取消订单)
*/
private void handleOrderMsg(WebSocketMessage wsMessage) {
/*
String senderId = wsMessage.getSenderId();
String data = wsMessage.getData();
// 1. 解析订单业务数据
OrderMsgDTO orderMsg = JSON.parseObject(data, OrderMsgDTO.class);
String orderId = orderMsg.getOrderId();
String action = orderMsg.getAction(); // 如QUERY/CANCEL
// 2. 执行业务逻辑
String result;
if ("QUERY".equals(action)) {
result = orderService.queryOrderById(orderId, senderId);
} else if ("CANCEL".equals(action)) {
result = orderService.cancelOrderById(orderId, senderId);
} else {
result = "不支持的订单操作:" + action;
}
// 3. 将业务结果推送给发送者
ChatWebSocketServer.sendToUser(senderId, "【订单处理结果】" + result);
*/
}
/**
* 处理心跳包(仅返回确认,无需复杂逻辑)
*/
private void handleHeartbeatMsg(WebSocketMessage wsMessage) {
String senderId = wsMessage.getSenderId();
// 回复心跳确认
ChatWebSocketServer.sendToUser(senderId, JSON.toJSONString(new WebSocketMessage("HEARTBEAT_ACK", "pong", senderId, null, System.currentTimeMillis())));
}
}

View File

@ -0,0 +1,78 @@
package com.czlis.interfaceCommon.websocket;
import javax.websocket.Session;
import java.net.InetSocketAddress;
import java.time.LocalDateTime;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* 封装WebSocket用户连接信息:聚合Session + 身份属性 + 连接元数据
*/
public class WebSocketUser {
// 核心字段:用户唯一标识
private final String userId;
// 身份属性(可根据业务扩展)
private final String username; // 用户名
private final String role; // 角色(如ADMIN/USER)
private final String loginIp; // 登录IP
// 连接元数据
private final Session session; // WebSocket连接会话
private final LocalDateTime connectTime; // 连接建立时间
// 动态属性(支持运行时更新,线程安全)
private final Map<String, Object> dynamicAttrs; // 自定义动态属性
// 构造器:初始化核心属性(必填)
public WebSocketUser(String userId, String username, String role, String loginIp, Session session) {
this.userId = userId;
this.username = username;
this.role = role;
this.loginIp = loginIp;
this.session = session;
this.connectTime = LocalDateTime.now();
this.dynamicAttrs = new ConcurrentHashMap<>(); // 动态属性用线程安全Map
}
// ========== 只读属性(getter) ==========
public String getUserId() {
return userId;
}
public String getUsername() {
return username;
}
public String getRole() {
return role;
}
public String getLoginIp() {
return loginIp;
}
public Session getSession() {
return session;
}
public LocalDateTime getConnectTime() {
return connectTime;
}
// ========== 动态属性(支持增删改查,线程安全) ==========
public void setDynamicAttr(String key, Object value) {
dynamicAttrs.put(key, value);
}
public Object getDynamicAttr(String key) {
return dynamicAttrs.get(key);
}
public void removeDynamicAttr(String key) {
dynamicAttrs.remove(key);
}
// 检查连接是否有效
public boolean isConnected() {
return session != null && session.isOpen();
}
}

View File

@ -0,0 +1,57 @@
websocket消息服务器结构说明
1.核心框架
ChatWebSocketServer
主要作用:提供底层的客户端连接管理,消息的发送与接收,以及客户端身份信息的修改
WebSocketMessage
说明:消息结构标准实体类
WebSocketUser
说明:客户端列表属性实体类
WebSocketMsgService
说明:WebSocket消息处理业务接口
WebSocketMsgServiceImpl
说明:WebSocket消息处理业务实现类
RetryMsgService
说明:消息重试服务,发送失败重发,最多3次
RetryMsgTask
说明:定时器服务器,增加定时业务服务的地方
2,协议约定
消息服务单实例,避免重复实例化错误
消息服务器地址/ws/msgserver/{SID}
SID为客户端唯一ID,约定使用浏览器指纹
WebSocketMessage消息结构
// 消息类型(如CHAT/NOTICE/ORDER/HEARTBEAT)
private String msgType;
// 业务数据(JSON字符串,可解析为对应业务DTO)
private String data;
// 发送者ID(可选,也可从WebSocket连接中获取)
private String senderId;
// 接收者ID(可选,单聊/广播区分)
private String receiverId;
// 消息时间戳
private Long timestamp;
消息时间戳必传,目前代码直接默认了
重点字段msgType,申明发送的是什么消息
发送者ID不需要传,自动获取
接收者ID在指定传输目标时传入,目前暂时不使用
data 字段为消息体,都是json格式
消息举例
{"msgType":"CHAT","data":{"instrid":"abc"},"receiverId":"LIS","timestamp":1755068400000}
{"msgType":"CHAT","data":{"instrid":"abc"}}
消息类型约定
3 业务流程
启动服务器,自动创建websocket服务实例
客户端启动,连接websocket服务器,发送约定身份参数(非必须)
客户端等待消息中...
服务器创建消息服务业务层
定时器注册指定服务为定时任务,并声明循环周期(单线程服务,一类业务只创建一个消息服务器端)
消息服务程序触发消息规则,给指定客户端发送消息
客户端收到消息,做相应的处理
4 应用范例
UserService
提供基本的消息服务方法,
常用的操作客户端方法直接操作ChatWebSocketServer完成
包括调阅客户端列表,修改客户端附属属性,过滤,查找,给目标客户端发送消息,以及处理收到的消息