Commit 6dabbff2 by zhangxingmin

push

parent e8a82078
......@@ -13,9 +13,10 @@ import com.yd.communication.service.service.ICoOperationLogService;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.PostConstruct;
import javax.websocket.*;
import javax.websocket.server.PathParam;
import javax.websocket.server.ServerEndpoint;
......@@ -155,46 +156,84 @@ public class CoWebSocketServer {
// ==================== 跨节点订阅(每个节点启动时执行) ====================
/**
* Bean 初始化完成后执行,启动 Redis 订阅线程
* 每个节点启动时都会订阅通配符频道 "room:channel:*"
* 用于接收其他节点发来的跨节点广播消息
* 监听应用启动完成事件,在容器完全就绪后启动 Redis 订阅
* 替代 @PostConstruct,避免 Bean 未完全初始化时连接 Redis
*/
@PostConstruct
public void init() {
// 启动一个独立的后台线程,避免阻塞主线程
@EventListener(ApplicationReadyEvent.class)
public void startRedisSubscriber() {
// 启动后台线程订阅 Redis
new Thread(() -> {
// 无限循环,支持断线自动重连
while (true) {
log.info("应用已就绪,开始启动 Redis 跨节点订阅...");
doSubscribe();
}).start();
}
/**
* 执行 Redis 订阅(支持自动重连)
*/
private void doSubscribe() {
while (true) {
try {
redisTemplate.execute((connection) -> {
connection.subscribe(
(message, pattern) -> handleCrossNodeMessage(new String(message.getBody())),
"room:channel:*".getBytes()
);
return null;
}, true);
} catch (Exception e) {
log.error("Redis 订阅断开,5秒后重试...", e);
try {
// 通过 Redis 连接执行订阅命令
redisTemplate.execute((connection) -> {
// 订阅所有以 "room:channel:" 开头的频道
// 第二个参数是监听器:收到消息时调用 handleCrossNodeMessage 处理
connection.subscribe(
(message, pattern) -> handleCrossNodeMessage(new String(message.getBody())),
"room:channel:*".getBytes()
);
// 返回 null,因为 subscribe 是阻塞方法,执行到这里说明订阅已结束(异常断开)
return null;
}, true); // true 表示使用事务(此处无实际影响)
} catch (Exception e) {
// 订阅断开(如 Redis 连接超时、网络抖动),记录错误日志
log.error("Redis 订阅断开,5秒后重试...", e);
try {
// 等待 5 秒后重连
Thread.sleep(5000);
} catch (InterruptedException ex) {
// 线程被中断,退出循环
Thread.currentThread().interrupt();
break;
}
Thread.sleep(5000);
} catch (InterruptedException ex) {
Thread.currentThread().interrupt();
break;
}
}
}).start(); // 启动线程
log.info("Redis 跨节点广播订阅已启动,当前节点ID: {}", NODE_ID);
}
}
/**
* Bean 初始化完成后执行,启动 Redis 订阅线程
* 每个节点启动时都会订阅通配符频道 "room:channel:*"
* 用于接收其他节点发来的跨节点广播消息
*/
// @PostConstruct
// public void init() {
// // 启动一个独立的后台线程,避免阻塞主线程
// new Thread(() -> {
// // 无限循环,支持断线自动重连
// while (true) {
// try {
// // 通过 Redis 连接执行订阅命令
// redisTemplate.execute((connection) -> {
// // 订阅所有以 "room:channel:" 开头的频道
// // 第二个参数是监听器:收到消息时调用 handleCrossNodeMessage 处理
// connection.subscribe(
// (message, pattern) -> handleCrossNodeMessage(new String(message.getBody())),
// "room:channel:*".getBytes()
// );
// // 返回 null,因为 subscribe 是阻塞方法,执行到这里说明订阅已结束(异常断开)
// return null;
// }, true); // true 表示使用事务(此处无实际影响)
// } catch (Exception e) {
// // 订阅断开(如 Redis 连接超时、网络抖动),记录错误日志
// log.error("Redis 订阅断开,5秒后重试...", e);
// try {
// // 等待 5 秒后重连
// Thread.sleep(5000);
// } catch (InterruptedException ex) {
// // 线程被中断,退出循环
// Thread.currentThread().interrupt();
// break;
// }
// }
// }
// }).start(); // 启动线程
// log.info("Redis 跨节点广播订阅已启动,当前节点ID: {}", NODE_ID);
// }
/**
* 处理跨节点广播消息(由 Redis 订阅触发)
* 当其他节点向 Redis 频道发布消息时,此方法会被调用
* 注意:此方法仅负责将消息转发给本节点的 Session,不处理状态同步
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment