Compare commits
2 Commits
6f0262ad04
...
805d9a2ffa
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
805d9a2ffa | ||
|
|
dbea5a7b3f |
@ -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;
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
@ -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();
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -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();
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -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;
|
||||||
|
}
|
||||||
@ -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);
|
||||||
|
}
|
||||||
@ -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())));
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -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();
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -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完成
|
||||||
|
包括调阅客户端列表,修改客户端附属属性,过滤,查找,给目标客户端发送消息,以及处理收到的消息
|
||||||
Loading…
x
Reference in New Issue
Block a user