Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
Y
yd-communication
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
xingmin
yd-communication
Commits
61ddf07a
Commit
61ddf07a
authored
Jul 30, 2026
by
zhangxingmin
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
push
parent
a9fd76ea
Hide whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
202 additions
and
81 deletions
+202
-81
yd-communication-api/src/main/java/com/yd/communication/api/websocket/CoWebSocketServer.java
+202
-81
No files found.
yd-communication-api/src/main/java/com/yd/communication/api/websocket/CoWebSocketServer.java
View file @
61ddf07a
...
...
@@ -287,103 +287,224 @@ public class CoWebSocketServer {
*/
private
final
ObjectMapper
objectMapper
=
new
ObjectMapper
();
/**
* WebSocket 连接建立时触发
* 1. 解析 URL 查询参数(userId, userType)
* 2. 将当前会话加入本地内存(ROOMS, SESSION_ROOM 等)
* 3. 将成员信息存入 Redis(跨节点共享)
* 4. 从 Redis 拉取当前房间状态并同步给新连接的用户(状态同步)
*
* @param session 当前 WebSocket 会话对象
* @param roomId 路径参数:房间号
*/
@OnOpen
public
void
onOpen
(
Session
session
,
@PathParam
(
"roomId"
)
String
roomId
)
{
// 1. 获取 URL 查询字符串(如 ?userId=123&userType=owner)
log
.
info
(
"========== WebSocket 连接请求开始 =========="
);
log
.
info
(
"房间号: {}, 会话ID: {}"
,
roomId
,
session
.
getId
());
// 1. 解析 URL 查询参数
String
queryString
=
session
.
getQueryString
();
// 解析查询参数为 Map
log
.
info
(
"查询字符串: {}"
,
queryString
);
Map
<
String
,
String
>
params
=
parseQueryString
(
queryString
);
// 获取用户ID,默认为 "unknown"
String
userId
=
params
.
getOrDefault
(
"userId"
,
"unknown"
);
// 获取用户类型,默认为 "unknown"(owner: 客户, participant: 顾问)
String
userType
=
params
.
getOrDefault
(
"userType"
,
"unknown"
);
log
.
info
(
"解析到的 userId: {}, userType: {}"
,
userId
,
userType
);
// 2. 检查关键依赖是否注入成功
log
.
info
(
"检查依赖注入 - sessionService: {}, redisTemplate: {}, redissonClient: {}"
,
sessionService
==
null
?
"NULL"
:
"OK"
,
redisTemplate
==
null
?
"NULL"
:
"OK"
,
redissonClient
==
null
?
"NULL"
:
"OK"
);
if
(
sessionService
==
null
||
redisTemplate
==
null
||
redissonClient
==
null
)
{
log
.
error
(
"依赖注入失败,无法处理连接!"
);
try
{
session
.
close
(
new
CloseReason
(
CloseReason
.
CloseCodes
.
UNEXPECTED_CONDITION
,
"服务未就绪"
));
}
catch
(
IOException
e
)
{
log
.
error
(
"关闭session失败"
,
e
);
}
return
;
}
// 2. 将当前会话加入本地内存
// computeIfAbsent: 如果房间不存在则创建新的 Set 集合
ROOMS
.
computeIfAbsent
(
roomId
,
k
->
new
CopyOnWriteArraySet
<>()).
add
(
session
);
// 记录会话对应的房间号
SESSION_ROOM
.
put
(
session
,
roomId
);
// 记录会话对应的用户ID
SESSION_USER_ID
.
put
(
session
,
userId
);
// 记录会话对应的用户类型
SESSION_USER_TYPE
.
put
(
session
,
userType
);
// 3. 将成员信息存入 Redis(供跨节点查询使用)
// 构造 Redis Hash 的键名:room:members:{roomId}
String
memberKey
=
String
.
format
(
ROOM_MEMBERS_KEY
,
roomId
);
// 构造成员信息 Map
Map
<
String
,
String
>
memberInfo
=
new
HashMap
<>();
memberInfo
.
put
(
"userId"
,
userId
);
memberInfo
.
put
(
"userType"
,
userType
);
memberInfo
.
put
(
"sessionId"
,
session
.
getId
());
// sessionId 作为唯一标识
// 以 sessionId 为 field,成员信息 JSON 为 value,存入 Redis Hash
redisTemplate
.
opsForHash
().
put
(
memberKey
,
session
.
getId
(),
JSON
.
toJSONString
(
memberInfo
));
// 设置 Hash 过期时间为 60 分钟(与会话超时一致)
redisTemplate
.
expire
(
memberKey
,
60
,
TimeUnit
.
MINUTES
);
// 4. 状态同步(核心功能):向新用户推送当前房间的完整状态
try
{
// 4.1 获取当前控制权信息(从 Redis 缓存中读取)
RoomRedisInfoDTO
dto
=
sessionService
.
getCurrentController
(
roomId
);
// 控制权持有者类型:owner 或 participant
String
controlHolderType
=
dto
.
getControlHolderType
();
// 控制权持有者ID
String
controlHolderId
=
dto
.
getControlHolderId
();
// 4.2 获取当前页面状态(从 Redis 中读取,由翻页操作实时更新)
String
stateKey
=
String
.
format
(
ROOM_STATE_PAGE_KEY
,
roomId
);
String
currentPageJson
=
redisTemplate
.
opsForValue
().
get
(
stateKey
);
// 如果 Redis 中没有页面状态,则从数据库读取最新状态
if
(
StringUtils
.
isBlank
(
currentPageJson
))
{
// 根据房间号查询数据库中的会话记录
CoSession
coSession
=
sessionService
.
getByRoomId
(
roomId
);
// 如果会话存在,取其中的当前页面 JSON
currentPageJson
=
coSession
!=
null
?
coSession
.
getCurrentPage
()
:
"{}"
;
// 将数据库中的状态回填到 Redis,供后续新节点快速同步
if
(
StringUtils
.
isNotBlank
(
currentPageJson
))
{
redisTemplate
.
opsForValue
().
set
(
stateKey
,
currentPageJson
,
60
,
TimeUnit
.
MINUTES
);
// 3. 加入本地内存
log
.
info
(
"将会话加入本地内存 ROOMS..."
);
ROOMS
.
computeIfAbsent
(
roomId
,
k
->
new
CopyOnWriteArraySet
<>()).
add
(
session
);
SESSION_ROOM
.
put
(
session
,
roomId
);
SESSION_USER_ID
.
put
(
session
,
userId
);
SESSION_USER_TYPE
.
put
(
session
,
userType
);
log
.
info
(
"本地内存添加成功,当前房间 {} 的会话数: {}"
,
roomId
,
ROOMS
.
get
(
roomId
).
size
());
// 4. 将成员信息存入 Redis
String
memberKey
=
String
.
format
(
ROOM_MEMBERS_KEY
,
roomId
);
Map
<
String
,
String
>
memberInfo
=
new
HashMap
<>();
memberInfo
.
put
(
"userId"
,
userId
);
memberInfo
.
put
(
"userType"
,
userType
);
memberInfo
.
put
(
"sessionId"
,
session
.
getId
());
log
.
info
(
"准备将成员信息存入 Redis, key: {}, info: {}"
,
memberKey
,
memberInfo
);
try
{
redisTemplate
.
opsForHash
().
put
(
memberKey
,
session
.
getId
(),
JSON
.
toJSONString
(
memberInfo
));
redisTemplate
.
expire
(
memberKey
,
60
,
TimeUnit
.
MINUTES
);
log
.
info
(
"Redis 成员信息存储成功"
);
}
catch
(
Exception
e
)
{
log
.
error
(
"Redis 成员信息存储失败"
,
e
);
throw
e
;
// 继续抛出以便外层捕获
}
// 5. 状态同步
log
.
info
(
"开始状态同步..."
);
try
{
// 5.1 获取控制权
RoomRedisInfoDTO
dto
=
sessionService
.
getCurrentController
(
roomId
);
String
controlHolderType
=
dto
.
getControlHolderType
();
String
controlHolderId
=
dto
.
getControlHolderId
();
log
.
info
(
"获取控制权成功: holderType={}, holderId={}"
,
controlHolderType
,
controlHolderId
);
// 5.2 获取当前页面状态
String
stateKey
=
String
.
format
(
ROOM_STATE_PAGE_KEY
,
roomId
);
String
currentPageJson
=
redisTemplate
.
opsForValue
().
get
(
stateKey
);
log
.
info
(
"从Redis获取当前页面状态: {}"
,
currentPageJson
);
if
(
StringUtils
.
isBlank
(
currentPageJson
))
{
CoSession
coSession
=
sessionService
.
getByRoomId
(
roomId
);
currentPageJson
=
coSession
!=
null
?
coSession
.
getCurrentPage
()
:
"{}"
;
if
(
StringUtils
.
isNotBlank
(
currentPageJson
))
{
redisTemplate
.
opsForValue
().
set
(
stateKey
,
currentPageJson
,
60
,
TimeUnit
.
MINUTES
);
log
.
info
(
"从数据库获取并回填Redis状态: {}"
,
currentPageJson
);
}
else
{
log
.
warn
(
"未找到当前页面状态,使用默认空对象"
);
}
}
// 5.3 发送初始化消息
String
initMsg
=
String
.
format
(
"{\"action\":\"init_sync\",\"holderType\":\"%s\",\"holderId\":\"%s\",\"currentPage\":%s}"
,
controlHolderType
,
controlHolderId
,
currentPageJson
);
session
.
getBasicRemote
().
sendText
(
initMsg
);
log
.
info
(
"发送 init_sync 消息: {}"
,
initMsg
);
String
controlMsg
=
String
.
format
(
"{\"action\":\"control_transfer\",\"holderType\":\"%s\",\"holderId\":\"%s\"}"
,
controlHolderType
,
controlHolderId
);
session
.
getBasicRemote
().
sendText
(
controlMsg
);
log
.
info
(
"发送 control_transfer 消息: {}"
,
controlMsg
);
log
.
info
(
"用户 {} 加入房间 {},状态同步完成"
,
userId
,
roomId
);
}
catch
(
Exception
e
)
{
log
.
error
(
"状态同步过程中发生异常"
,
e
);
// 状态同步失败不影响连接建立,但应告知客户端(可选)
// 这里重新抛出以便外层统一处理
throw
e
;
}
// 4.3 组装初始化同步消息(包含控制权和当前页面)
// 消息格式:{"action":"init_sync","holderType":"xxx","holderId":"xxx","currentPage":{...}}
String
initMsg
=
String
.
format
(
"{\"action\":\"init_sync\",\"holderType\":\"%s\",\"holderId\":\"%s\",\"currentPage\":%s}"
,
controlHolderType
,
controlHolderId
,
currentPageJson
);
// 发送初始化同步消息给刚连接的客户端
session
.
getBasicRemote
().
sendText
(
initMsg
);
// 4.4 单独再推送一次控制权消息(兼容前端只监听 control_transfer 的情况)
String
controlMsg
=
String
.
format
(
"{\"action\":\"control_transfer\",\"holderType\":\"%s\",\"holderId\":\"%s\"}"
,
controlHolderType
,
controlHolderId
);
session
.
getBasicRemote
().
sendText
(
controlMsg
);
// 记录日志:用户已加入并完成状态同步
log
.
info
(
"用户 {} 加入房间 {},已同步状态"
,
userId
,
roomId
);
log
.
info
(
"========== WebSocket 连接处理完成 =========="
);
}
catch
(
Exception
e
)
{
// 状态同步失败,记录错误但不影响连接建立
log
.
error
(
"状态同步失败,用户 {} 可能无法恢复最新状态"
,
userId
,
e
);
log
.
error
(
"WebSocket 连接处理失败,房间号: {}, 用户: {}"
,
roomId
,
userId
,
e
);
// 发送错误消息给客户端(可选)
try
{
session
.
getBasicRemote
().
sendText
(
"{\"error\":\"服务器处理异常: "
+
e
.
getMessage
()
+
"\"}"
);
}
catch
(
IOException
ex
)
{
log
.
error
(
"发送错误消息失败"
,
ex
);
}
// 关闭连接
try
{
session
.
close
(
new
CloseReason
(
CloseReason
.
CloseCodes
.
UNEXPECTED_CONDITION
,
"服务器错误: "
+
e
.
getMessage
()));
}
catch
(
IOException
ex
)
{
log
.
error
(
"关闭session失败"
,
ex
);
}
}
// 记录连接建立日志,包含当前节点该房间的会话数
log
.
info
(
"用户 {} ({}) 加入房间 {}"
,
userId
,
userType
,
roomId
);
}
/**
* WebSocket 连接建立时触发
* 1. 解析 URL 查询参数(userId, userType)
* 2. 将当前会话加入本地内存(ROOMS, SESSION_ROOM 等)
* 3. 将成员信息存入 Redis(跨节点共享)
* 4. 从 Redis 拉取当前房间状态并同步给新连接的用户(状态同步)
*
* @param session 当前 WebSocket 会话对象
* @param roomId 路径参数:房间号
*/
// @OnOpen
// public void onOpen(Session session, @PathParam("roomId") String roomId) {
// // 1. 获取 URL 查询字符串(如 ?userId=123&userType=owner)
// String queryString = session.getQueryString();
// // 解析查询参数为 Map
// Map<String, String> params = parseQueryString(queryString);
// // 获取用户ID,默认为 "unknown"
// String userId = params.getOrDefault("userId", "unknown");
// // 获取用户类型,默认为 "unknown"(owner: 客户, participant: 顾问)
// String userType = params.getOrDefault("userType", "unknown");
//
// // 2. 将当前会话加入本地内存
// // computeIfAbsent: 如果房间不存在则创建新的 Set 集合
// ROOMS.computeIfAbsent(roomId, k -> new CopyOnWriteArraySet<>()).add(session);
// // 记录会话对应的房间号
// SESSION_ROOM.put(session, roomId);
// // 记录会话对应的用户ID
// SESSION_USER_ID.put(session, userId);
// // 记录会话对应的用户类型
// SESSION_USER_TYPE.put(session, userType);
//
// // 3. 将成员信息存入 Redis(供跨节点查询使用)
// // 构造 Redis Hash 的键名:room:members:{roomId}
// String memberKey = String.format(ROOM_MEMBERS_KEY, roomId);
// // 构造成员信息 Map
// Map<String, String> memberInfo = new HashMap<>();
// memberInfo.put("userId", userId);
// memberInfo.put("userType", userType);
// memberInfo.put("sessionId", session.getId()); // sessionId 作为唯一标识
// // 以 sessionId 为 field,成员信息 JSON 为 value,存入 Redis Hash
// redisTemplate.opsForHash().put(memberKey, session.getId(), JSON.toJSONString(memberInfo));
// // 设置 Hash 过期时间为 60 分钟(与会话超时一致)
// redisTemplate.expire(memberKey, 60, TimeUnit.MINUTES);
//
// // 4. 状态同步(核心功能):向新用户推送当前房间的完整状态
// try {
// // 4.1 获取当前控制权信息(从 Redis 缓存中读取)
// RoomRedisInfoDTO dto = sessionService.getCurrentController(roomId);
// // 控制权持有者类型:owner 或 participant
// String controlHolderType = dto.getControlHolderType();
// // 控制权持有者ID
// String controlHolderId = dto.getControlHolderId();
//
// // 4.2 获取当前页面状态(从 Redis 中读取,由翻页操作实时更新)
// String stateKey = String.format(ROOM_STATE_PAGE_KEY, roomId);
// String currentPageJson = redisTemplate.opsForValue().get(stateKey);
//
// // 如果 Redis 中没有页面状态,则从数据库读取最新状态
// if (StringUtils.isBlank(currentPageJson)) {
// // 根据房间号查询数据库中的会话记录
// CoSession coSession = sessionService.getByRoomId(roomId);
// // 如果会话存在,取其中的当前页面 JSON
// currentPageJson = coSession != null ? coSession.getCurrentPage() : "{}";
// // 将数据库中的状态回填到 Redis,供后续新节点快速同步
// if (StringUtils.isNotBlank(currentPageJson)) {
// redisTemplate.opsForValue().set(stateKey, currentPageJson, 60, TimeUnit.MINUTES);
// }
// }
//
// // 4.3 组装初始化同步消息(包含控制权和当前页面)
// // 消息格式:{"action":"init_sync","holderType":"xxx","holderId":"xxx","currentPage":{...}}
// String initMsg = String.format(
// "{\"action\":\"init_sync\",\"holderType\":\"%s\",\"holderId\":\"%s\",\"currentPage\":%s}",
// controlHolderType, controlHolderId, currentPageJson
// );
// // 发送初始化同步消息给刚连接的客户端
// session.getBasicRemote().sendText(initMsg);
//
// // 4.4 单独再推送一次控制权消息(兼容前端只监听 control_transfer 的情况)
// String controlMsg = String.format(
// "{\"action\":\"control_transfer\",\"holderType\":\"%s\",\"holderId\":\"%s\"}",
// controlHolderType, controlHolderId
// );
// session.getBasicRemote().sendText(controlMsg);
//
// // 记录日志:用户已加入并完成状态同步
// log.info("用户 {} 加入房间 {},已同步状态", userId, roomId);
// } catch (Exception e) {
// // 状态同步失败,记录错误但不影响连接建立
// log.error("状态同步失败,用户 {} 可能无法恢复最新状态", userId, e);
// }
//
// // 记录连接建立日志,包含当前节点该房间的会话数
// log.info("用户 {} ({}) 加入房间 {}", userId, userType, roomId);
// }
/**
* 接收客户端发送的 WebSocket 消息(核心业务分发器)
* 根据 action 类型分发到不同的业务处理逻辑
*
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment