From bb3f9cc5e40eab66af078eec88f744a320ca35bd Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?=E9=82=93=E6=96=87=E5=85=B5?= <441889070@qq.com>
Date: Wed, 13 Aug 2025 10:08:25 +0800
Subject: [PATCH] =?UTF-8?q?=E9=A1=B9=E7=9B=AE=E8=BF=81=E7=A7=BB?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
.../.flattened-pom.xml | 92 +++++++++
.../pom.xml | 73 +++++++
.../websocket/config/WebSocketProperties.java | 34 ++++
.../YudaoWebSocketAutoConfiguration.java | 183 ++++++++++++++++++
.../handler/JsonWebSocketMessageHandler.java | 83 ++++++++
.../listener/WebSocketMessageListener.java | 31 +++
.../core/message/JsonWebSocketMessage.java | 29 +++
.../LoginUserHandshakeInterceptor.java | 42 ++++
.../WebSocketAuthorizeRequestsCustomizer.java | 24 +++
.../AbstractWebSocketMessageSender.java | 106 ++++++++++
.../core/sender/WebSocketMessageSender.java | 52 +++++
.../sender/kafka/KafkaWebSocketMessage.java | 35 ++++
.../kafka/KafkaWebSocketMessageConsumer.java | 28 +++
.../kafka/KafkaWebSocketMessageSender.java | 67 +++++++
.../local/LocalWebSocketMessageSender.java | 20 ++
.../rabbitmq/RabbitMQWebSocketMessage.java | 37 ++++
.../RabbitMQWebSocketMessageConsumer.java | 39 ++++
.../RabbitMQWebSocketMessageSender.java | 62 ++++++
.../sender/redis/RedisWebSocketMessage.java | 34 ++++
.../redis/RedisWebSocketMessageConsumer.java | 23 +++
.../redis/RedisWebSocketMessageSender.java | 57 ++++++
.../rocketmq/RocketMQWebSocketMessage.java | 35 ++++
.../RocketMQWebSocketMessageConsumer.java | 30 +++
.../RocketMQWebSocketMessageSender.java | 61 ++++++
.../WebSocketSessionHandlerDecorator.java | 49 +++++
.../core/session/WebSocketSessionManager.java | 53 +++++
.../session/WebSocketSessionManagerImpl.java | 125 ++++++++++++
.../core/util/WebSocketFrameworkUtils.java | 67 +++++++
.../framework/websocket/package-info.java | 4 +
...ot.autoconfigure.AutoConfiguration.imports | 1 +
.../《芋道 Spring Boot WebSocket 入门》.md | 1 +
31 files changed, 1577 insertions(+)
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/.flattened-pom.xml
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/pom.xml
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/config/WebSocketProperties.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/config/YudaoWebSocketAutoConfiguration.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/handler/JsonWebSocketMessageHandler.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/listener/WebSocketMessageListener.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/message/JsonWebSocketMessage.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/security/LoginUserHandshakeInterceptor.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/security/WebSocketAuthorizeRequestsCustomizer.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/AbstractWebSocketMessageSender.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/WebSocketMessageSender.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/kafka/KafkaWebSocketMessage.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/kafka/KafkaWebSocketMessageConsumer.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/kafka/KafkaWebSocketMessageSender.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/local/LocalWebSocketMessageSender.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/rabbitmq/RabbitMQWebSocketMessage.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/rabbitmq/RabbitMQWebSocketMessageConsumer.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/rabbitmq/RabbitMQWebSocketMessageSender.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/redis/RedisWebSocketMessage.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/redis/RedisWebSocketMessageConsumer.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/redis/RedisWebSocketMessageSender.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/rocketmq/RocketMQWebSocketMessage.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/rocketmq/RocketMQWebSocketMessageConsumer.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/sender/rocketmq/RocketMQWebSocketMessageSender.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/session/WebSocketSessionHandlerDecorator.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/session/WebSocketSessionManager.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/session/WebSocketSessionManagerImpl.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/core/util/WebSocketFrameworkUtils.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/package-info.java
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports
create mode 100644 yudao-framework/yudao-spring-boot-starter-websocket/《芋道 Spring Boot WebSocket 入门》.md
diff --git a/yudao-framework/yudao-spring-boot-starter-websocket/.flattened-pom.xml b/yudao-framework/yudao-spring-boot-starter-websocket/.flattened-pom.xml
new file mode 100644
index 00000000..bb525b8d
--- /dev/null
+++ b/yudao-framework/yudao-spring-boot-starter-websocket/.flattened-pom.xml
@@ -0,0 +1,92 @@
+
+
+ 4.0.0
+ cn.iocoder.boot
+ yudao-spring-boot-starter-websocket
+ 2025.08-SNAPSHOT
+ yudao-spring-boot-starter-websocket
+ WebSocket 框架,支持多节点的广播
+ https://github.com/YunaiV/ruoyi-vue-pro
+
+
+ cn.iocoder.boot
+ yudao-common
+ 2025.08-SNAPSHOT
+ compile
+
+
+ cn.iocoder.boot
+ yudao-spring-boot-starter-security
+ 2025.08-SNAPSHOT
+ provided
+
+
+ org.springframework.boot
+ spring-boot-starter-websocket
+ 3.4.5
+ compile
+
+
+ cn.iocoder.boot
+ yudao-spring-boot-starter-mq
+ 2025.08-SNAPSHOT
+ compile
+
+
+ org.springframework.kafka
+ spring-kafka
+ 3.3.5
+ compile
+ true
+
+
+ org.springframework.amqp
+ spring-rabbit
+ 3.2.5
+ compile
+ true
+
+
+ org.apache.rocketmq
+ rocketmq-spring-boot-starter
+ 2.3.2
+ compile
+ true
+
+
+ cn.iocoder.boot
+ yudao-spring-boot-starter-biz-tenant
+ 2025.08-SNAPSHOT
+ provided
+
+
+
+
+ huaweicloud
+ huawei
+ https://mirrors.huaweicloud.com/repository/maven/
+
+
+ aliyunmaven
+ aliyun
+ https://maven.aliyun.com/repository/public
+
+
+
+ false
+
+ spring-milestones
+ Spring Milestones
+ https://repo.spring.io/milestone
+
+
+
+ false
+
+ spring-snapshots
+ Spring Snapshots
+ https://repo.spring.io/snapshot
+
+
+
diff --git a/yudao-framework/yudao-spring-boot-starter-websocket/pom.xml b/yudao-framework/yudao-spring-boot-starter-websocket/pom.xml
new file mode 100644
index 00000000..b534f10d
--- /dev/null
+++ b/yudao-framework/yudao-spring-boot-starter-websocket/pom.xml
@@ -0,0 +1,73 @@
+
+
+
+ cn.iocoder.boot
+ yudao-framework
+ ${revision}
+
+ 4.0.0
+ yudao-spring-boot-starter-websocket
+ jar
+
+ ${project.artifactId}
+ WebSocket 框架,支持多节点的广播
+ https://github.com/YunaiV/ruoyi-vue-pro
+
+
+
+
+ cn.iocoder.boot
+ yudao-common
+
+
+
+
+
+ cn.iocoder.boot
+ yudao-spring-boot-starter-security
+ provided
+
+
+
+ org.springframework.boot
+ spring-boot-starter-websocket
+
+
+
+
+ cn.iocoder.boot
+ yudao-spring-boot-starter-mq
+
+
+ org.springframework.kafka
+ spring-kafka
+ true
+
+
+ org.springframework.amqp
+ spring-rabbit
+ true
+
+
+ org.apache.rocketmq
+ rocketmq-spring-boot-starter
+ true
+
+
+
+
+
+ cn.iocoder.boot
+ yudao-spring-boot-starter-biz-tenant
+ provided
+
+
+
+
\ No newline at end of file
diff --git a/yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/config/WebSocketProperties.java b/yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/config/WebSocketProperties.java
new file mode 100644
index 00000000..f1d84f7a
--- /dev/null
+++ b/yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/config/WebSocketProperties.java
@@ -0,0 +1,34 @@
+package cn.iocoder.yudao.framework.websocket.config;
+
+import lombok.Data;
+import org.springframework.boot.context.properties.ConfigurationProperties;
+import org.springframework.validation.annotation.Validated;
+
+import jakarta.validation.constraints.NotEmpty;
+import jakarta.validation.constraints.NotNull;
+
+/**
+ * WebSocket 配置项
+ *
+ * @author xingyu4j
+ */
+@ConfigurationProperties("yudao.websocket")
+@Data
+@Validated
+public class WebSocketProperties {
+
+ /**
+ * WebSocket 的连接路径
+ */
+ @NotEmpty(message = "WebSocket 的连接路径不能为空")
+ private String path = "/ws";
+
+ /**
+ * 消息发送器的类型
+ *
+ * 可选值:local、redis、rocketmq、kafka、rabbitmq
+ */
+ @NotNull(message = "WebSocket 的消息发送者不能为空")
+ private String senderType = "local";
+
+}
diff --git a/yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/config/YudaoWebSocketAutoConfiguration.java b/yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/config/YudaoWebSocketAutoConfiguration.java
new file mode 100644
index 00000000..3aded887
--- /dev/null
+++ b/yudao-framework/yudao-spring-boot-starter-websocket/src/main/java/cn/iocoder/yudao/framework/websocket/config/YudaoWebSocketAutoConfiguration.java
@@ -0,0 +1,183 @@
+package cn.iocoder.yudao.framework.websocket.config;
+
+import cn.iocoder.yudao.framework.mq.redis.config.YudaoRedisMQConsumerAutoConfiguration;
+import cn.iocoder.yudao.framework.mq.redis.core.RedisMQTemplate;
+import cn.iocoder.yudao.framework.websocket.core.handler.JsonWebSocketMessageHandler;
+import cn.iocoder.yudao.framework.websocket.core.listener.WebSocketMessageListener;
+import cn.iocoder.yudao.framework.websocket.core.security.LoginUserHandshakeInterceptor;
+import cn.iocoder.yudao.framework.websocket.core.security.WebSocketAuthorizeRequestsCustomizer;
+import cn.iocoder.yudao.framework.websocket.core.sender.kafka.KafkaWebSocketMessageConsumer;
+import cn.iocoder.yudao.framework.websocket.core.sender.kafka.KafkaWebSocketMessageSender;
+import cn.iocoder.yudao.framework.websocket.core.sender.local.LocalWebSocketMessageSender;
+import cn.iocoder.yudao.framework.websocket.core.sender.rabbitmq.RabbitMQWebSocketMessageConsumer;
+import cn.iocoder.yudao.framework.websocket.core.sender.rabbitmq.RabbitMQWebSocketMessageSender;
+import cn.iocoder.yudao.framework.websocket.core.sender.redis.RedisWebSocketMessageConsumer;
+import cn.iocoder.yudao.framework.websocket.core.sender.redis.RedisWebSocketMessageSender;
+import cn.iocoder.yudao.framework.websocket.core.sender.rocketmq.RocketMQWebSocketMessageConsumer;
+import cn.iocoder.yudao.framework.websocket.core.sender.rocketmq.RocketMQWebSocketMessageSender;
+import cn.iocoder.yudao.framework.websocket.core.session.WebSocketSessionHandlerDecorator;
+import cn.iocoder.yudao.framework.websocket.core.session.WebSocketSessionManager;
+import cn.iocoder.yudao.framework.websocket.core.session.WebSocketSessionManagerImpl;
+import org.apache.rocketmq.spring.core.RocketMQTemplate;
+import org.springframework.amqp.core.TopicExchange;
+import org.springframework.amqp.rabbit.core.RabbitTemplate;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.boot.autoconfigure.AutoConfiguration;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.boot.context.properties.EnableConfigurationProperties;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.web.socket.WebSocketHandler;
+import org.springframework.web.socket.config.annotation.EnableWebSocket;
+import org.springframework.web.socket.config.annotation.WebSocketConfigurer;
+import org.springframework.web.socket.server.HandshakeInterceptor;
+
+import java.util.List;
+
+/**
+ * WebSocket 自动配置
+ *
+ * @author xingyu4j
+ */
+@AutoConfiguration(before = YudaoRedisMQConsumerAutoConfiguration.class) // before YudaoRedisMQConsumerAutoConfiguration 的原因是,需要保证 RedisWebSocketMessageConsumer 先创建,才能创建 RedisMessageListenerContainer
+@EnableWebSocket // 开启 websocket
+@ConditionalOnProperty(prefix = "yudao.websocket", value = "enable", matchIfMissing = true) // 允许使用 yudao.websocket.enable=false 禁用 websocket
+@EnableConfigurationProperties(WebSocketProperties.class)
+public class YudaoWebSocketAutoConfiguration {
+
+ @Bean
+ public WebSocketConfigurer webSocketConfigurer(HandshakeInterceptor[] handshakeInterceptors,
+ WebSocketHandler webSocketHandler,
+ WebSocketProperties webSocketProperties) {
+ return registry -> registry
+ // 添加 WebSocketHandler
+ .addHandler(webSocketHandler, webSocketProperties.getPath())
+ .addInterceptors(handshakeInterceptors)
+ // 允许跨域,否则前端连接会直接断开
+ .setAllowedOriginPatterns("*");
+ }
+
+ @Bean
+ public HandshakeInterceptor handshakeInterceptor() {
+ return new LoginUserHandshakeInterceptor();
+ }
+
+ @Bean
+ public WebSocketHandler webSocketHandler(WebSocketSessionManager sessionManager,
+ List extends WebSocketMessageListener>> messageListeners) {
+ // 1. 创建 JsonWebSocketMessageHandler 对象,处理消息
+ JsonWebSocketMessageHandler messageHandler = new JsonWebSocketMessageHandler(messageListeners);
+ // 2. 创建 WebSocketSessionHandlerDecorator 对象,处理连接
+ return new WebSocketSessionHandlerDecorator(messageHandler, sessionManager);
+ }
+
+ @Bean
+ public WebSocketSessionManager webSocketSessionManager() {
+ return new WebSocketSessionManagerImpl();
+ }
+
+ @Bean
+ public WebSocketAuthorizeRequestsCustomizer webSocketAuthorizeRequestsCustomizer(WebSocketProperties webSocketProperties) {
+ return new WebSocketAuthorizeRequestsCustomizer(webSocketProperties);
+ }
+
+ // ==================== Sender 相关 ====================
+
+ @Configuration
+ @ConditionalOnProperty(prefix = "yudao.websocket", name = "sender-type", havingValue = "local")
+ public class LocalWebSocketMessageSenderConfiguration {
+
+ @Bean
+ public LocalWebSocketMessageSender localWebSocketMessageSender(WebSocketSessionManager sessionManager) {
+ return new LocalWebSocketMessageSender(sessionManager);
+ }
+
+ }
+
+ @Configuration
+ @ConditionalOnProperty(prefix = "yudao.websocket", name = "sender-type", havingValue = "redis")
+ public class RedisWebSocketMessageSenderConfiguration {
+
+ @Bean
+ public RedisWebSocketMessageSender redisWebSocketMessageSender(WebSocketSessionManager sessionManager,
+ RedisMQTemplate redisMQTemplate) {
+ return new RedisWebSocketMessageSender(sessionManager, redisMQTemplate);
+ }
+
+ @Bean
+ public RedisWebSocketMessageConsumer redisWebSocketMessageConsumer(
+ RedisWebSocketMessageSender redisWebSocketMessageSender) {
+ return new RedisWebSocketMessageConsumer(redisWebSocketMessageSender);
+ }
+
+ }
+
+ @Configuration
+ @ConditionalOnProperty(prefix = "yudao.websocket", name = "sender-type", havingValue = "rocketmq")
+ public class RocketMQWebSocketMessageSenderConfiguration {
+
+ @Bean
+ public RocketMQWebSocketMessageSender rocketMQWebSocketMessageSender(
+ WebSocketSessionManager sessionManager, RocketMQTemplate rocketMQTemplate,
+ @Value("${yudao.websocket.sender-rocketmq.topic}") String topic) {
+ return new RocketMQWebSocketMessageSender(sessionManager, rocketMQTemplate, topic);
+ }
+
+ @Bean
+ public RocketMQWebSocketMessageConsumer rocketMQWebSocketMessageConsumer(
+ RocketMQWebSocketMessageSender rocketMQWebSocketMessageSender) {
+ return new RocketMQWebSocketMessageConsumer(rocketMQWebSocketMessageSender);
+ }
+
+ }
+
+ @Configuration
+ @ConditionalOnProperty(prefix = "yudao.websocket", name = "sender-type", havingValue = "rabbitmq")
+ public class RabbitMQWebSocketMessageSenderConfiguration {
+
+ @Bean
+ public RabbitMQWebSocketMessageSender rabbitMQWebSocketMessageSender(
+ WebSocketSessionManager sessionManager, RabbitTemplate rabbitTemplate,
+ TopicExchange websocketTopicExchange) {
+ return new RabbitMQWebSocketMessageSender(sessionManager, rabbitTemplate, websocketTopicExchange);
+ }
+
+ @Bean
+ public RabbitMQWebSocketMessageConsumer rabbitMQWebSocketMessageConsumer(
+ RabbitMQWebSocketMessageSender rabbitMQWebSocketMessageSender) {
+ return new RabbitMQWebSocketMessageConsumer(rabbitMQWebSocketMessageSender);
+ }
+
+ /**
+ * 创建 Topic Exchange
+ */
+ @Bean
+ public TopicExchange websocketTopicExchange(@Value("${yudao.websocket.sender-rabbitmq.exchange}") String exchange) {
+ return new TopicExchange(exchange,
+ true, // durable: 是否持久化
+ false); // exclusive: 是否排它
+ }
+
+ }
+
+ @Configuration
+ @ConditionalOnProperty(prefix = "yudao.websocket", name = "sender-type", havingValue = "kafka")
+ public class KafkaWebSocketMessageSenderConfiguration {
+
+ @Bean
+ public KafkaWebSocketMessageSender kafkaWebSocketMessageSender(
+ WebSocketSessionManager sessionManager, KafkaTemplate