From fa2f75b6865bdb87d2b4160763f0277d4f0d334a Mon Sep 17 00:00:00 2001 From: AprilWind <2100166581@qq.com> Date: Wed, 8 Jul 2026 10:41:24 +0800 Subject: [PATCH] =?UTF-8?q?update=20=E6=B7=BB=E5=8A=A0=20MQTT=20=E6=B6=88?= =?UTF-8?q?=E6=81=AF=E6=8E=A8=E9=80=81=E6=94=AF=E6=8C=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/main/resources/application.yml | 4 +- ruoyi-common/ruoyi-common-push/pom.xml | 5 + .../push/config/MessageMqttConfiguration.java | 34 +++++ .../push/constant/MessageConstants.java | 10 ++ .../push/core/MqttPushSessionManager.java | 120 ++++++++++++++++++ .../push/enums/MessageTransportEnum.java | 8 +- .../push/properties/MessageProperties.java | 8 +- ...ot.autoconfigure.AutoConfiguration.imports | 1 + 8 files changed, 187 insertions(+), 3 deletions(-) create mode 100644 ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/config/MessageMqttConfiguration.java create mode 100644 ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/core/MqttPushSessionManager.java diff --git a/ruoyi-admin/src/main/resources/application.yml b/ruoyi-admin/src/main/resources/application.yml index 32b8a7429..3a07c5cb9 100644 --- a/ruoyi-admin/src/main/resources/application.yml +++ b/ruoyi-admin/src/main/resources/application.yml @@ -217,7 +217,7 @@ management: --- # 统一消息推送配置 message: enabled: true - # sse / websocket + # sse / websocket / mqtt transport: sse # 统一访问路径 path: /resource/message @@ -231,6 +231,8 @@ message: web-socket-send-time-limit: 10000 # WebSocket 发送缓冲区大小 web-socket-buffer-size-limit: 64000 + # MQTT 消息推送主题前缀,广播主题: ${prefix}/all,用户主题: ${prefix}/user/{userId} + mqtt-topic-prefix: /message --- # warm-flow工作流配置 warm-flow: diff --git a/ruoyi-common/ruoyi-common-push/pom.xml b/ruoyi-common/ruoyi-common-push/pom.xml index f86a27f3a..98d143e0a 100644 --- a/ruoyi-common/ruoyi-common-push/pom.xml +++ b/ruoyi-common/ruoyi-common-push/pom.xml @@ -52,5 +52,10 @@ + + + org.dromara + ruoyi-common-mqtt + diff --git a/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/config/MessageMqttConfiguration.java b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/config/MessageMqttConfiguration.java new file mode 100644 index 000000000..d3ba9e35e --- /dev/null +++ b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/config/MessageMqttConfiguration.java @@ -0,0 +1,34 @@ +package org.dromara.common.push.config; + +import org.dromara.common.mqtt.config.MqttAutoConfiguration; +import org.dromara.common.push.annotation.ConditionalOnMessageTransport; +import org.dromara.common.push.core.MqttPushSessionManager; +import org.dromara.common.push.properties.MessageProperties; +import org.dromara.mica.mqtt.spring.client.MqttClientTemplate; +import org.springframework.boot.autoconfigure.AutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.context.annotation.Bean; + +/** + * MQTT 消息推送自动装配。 + * + * @author Lion Li + */ +@AutoConfiguration(after = {MessageAutoConfiguration.class, MqttAutoConfiguration.class}) +@ConditionalOnMessageTransport("mqtt") +@ConditionalOnProperty(value = "mqtt.client.enabled", havingValue = "true") +public class MessageMqttConfiguration { + + /** + * MQTT 推送会话管理器。 + * + * @param client MQTT 客户端模板 + * @param messageProperties 消息推送配置 + * @return MqttPushSessionManager 实例 + */ + @Bean + public MqttPushSessionManager mqttPushSessionManager(MqttClientTemplate client, + MessageProperties messageProperties) { + return new MqttPushSessionManager(client, messageProperties); + } +} diff --git a/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/constant/MessageConstants.java b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/constant/MessageConstants.java index f35a4a273..ed7354b1f 100644 --- a/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/constant/MessageConstants.java +++ b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/constant/MessageConstants.java @@ -22,6 +22,16 @@ public interface MessageConstants { */ String MESSAGE_TOPIC = "global:message"; + /** + * MQTT 广播消息主题后缀。 + */ + String MQTT_BROADCAST_TOPIC_SUFFIX = "/all"; + + /** + * MQTT 指定用户消息主题目录。 + */ + String MQTT_USER_TOPIC_PATH = "/user/"; + /** * 心跳请求标识 */ diff --git a/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/core/MqttPushSessionManager.java b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/core/MqttPushSessionManager.java new file mode 100644 index 000000000..70a5f755d --- /dev/null +++ b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/core/MqttPushSessionManager.java @@ -0,0 +1,120 @@ +package org.dromara.common.push.core; + +import cn.hutool.core.collection.CollUtil; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.dromara.common.json.utils.JsonUtils; +import org.dromara.common.push.constant.MessageConstants; +import org.dromara.common.push.dto.PushDTO; +import org.dromara.common.push.properties.MessageProperties; +import org.dromara.mica.mqtt.spring.client.MqttClientTemplate; +import org.dromara.system.api.domain.PushPayloadDTO; + +import java.nio.charset.StandardCharsets; +import java.util.function.Consumer; + +/** + * MQTT 推送会话管理器。 + * + * @author Lion Li + */ +@Slf4j +@RequiredArgsConstructor +public class MqttPushSessionManager implements PushSessionManager { + + private final MqttClientTemplate client; + + private final MessageProperties messageProperties; + + /** + * MQTT 推送不依赖 Redis 主题订阅。 + * + * @param consumer 消息消费逻辑 + */ + @Override + public void subscribeMessage(Consumer consumer) { + // MQTT 通道直接发布到 broker,不需要订阅 Redis 分发。 + } + + /** + * 发布指定用户 MQTT 消息。 + * + * @param userId 目标用户ID + * @param payload 消息体 + */ + @Override + public void sendMessage(Long userId, PushPayloadDTO payload) { + if (userId == null || payload == null) { + return; + } + publish(topicPrefix() + MessageConstants.MQTT_USER_TOPIC_PATH + userId, payload); + } + + /** + * 发布 MQTT 广播消息。 + * + * @param payload 消息体 + */ + @Override + public void sendMessage(PushPayloadDTO payload) { + if (payload == null) { + return; + } + publish(topicPrefix() + MessageConstants.MQTT_BROADCAST_TOPIC_SUFFIX, payload); + } + + /** + * 发布 MQTT 消息。 + * + * @param pushDTO 推送参数封装对象 + */ + @Override + public void publishMessage(PushDTO pushDTO) { + if (pushDTO == null || pushDTO.getPayload() == null) { + return; + } + if (CollUtil.isEmpty(pushDTO.getUserIds())) { + sendMessage(pushDTO.getPayload()); + return; + } + pushDTO.getUserIds().forEach(userId -> sendMessage(userId, pushDTO.getPayload())); + } + + /** + * 发布 MQTT 广播消息。 + * + * @param payload 消息体 + */ + @Override + public void publishAll(PushPayloadDTO payload) { + sendMessage(payload); + } + + /** + * 发送消息到 MQTT Broker。 + * + * @param topic MQTT 主题 + * @param payload 消息体 + */ + private void publish(String topic, PushPayloadDTO payload) { + String message = JsonUtils.toJsonString(payload); + client.publish(topic, message.getBytes(StandardCharsets.UTF_8)); + log.info("MQTT发送推送消息topic:{} message:{}", topic, payload.getMessage()); + } + + /** + * 获取规范化后的 MQTT topic 前缀。 + * + * @return MQTT topic 前缀 + */ + private String topicPrefix() { + String prefix = messageProperties.getMqttTopicPrefix(); + if (prefix == null || prefix.isBlank()) { + return "/message"; + } + while (prefix.endsWith("/") && prefix.length() > 1) { + prefix = prefix.substring(0, prefix.length() - 1); + } + return prefix; + } +} diff --git a/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/enums/MessageTransportEnum.java b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/enums/MessageTransportEnum.java index 362d26645..0efc3d071 100644 --- a/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/enums/MessageTransportEnum.java +++ b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/enums/MessageTransportEnum.java @@ -24,7 +24,13 @@ public enum MessageTransportEnum { * WebSocket 传输方式 * 全双工长连接,支持双向实时通信 */ - WEBSOCKET("websocket"); + WEBSOCKET("websocket"), + + /** + * MQTT 传输方式 + * 通过 MQTT Broker 发布统一推送消息 + */ + MQTT("mqtt"); /** * 传输类型编码 diff --git a/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/properties/MessageProperties.java b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/properties/MessageProperties.java index 7d0aeaebf..e20ccd996 100644 --- a/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/properties/MessageProperties.java +++ b/ruoyi-common/ruoyi-common-push/src/main/java/org/dromara/common/push/properties/MessageProperties.java @@ -19,7 +19,7 @@ public class MessageProperties { private Boolean enabled = true; /** - * 传输方式:sse / websocket。 + * 传输方式:sse / websocket / mqtt。 */ private String transport = MessageTransportEnum.SSE.getCode(); @@ -52,4 +52,10 @@ public class MessageProperties { * WebSocket 发送缓冲区大小。 */ private int webSocketBufferSizeLimit = 64_000; + + /** + * MQTT 消息推送主题前缀。 + */ + private String mqttTopicPrefix = "/message"; + } diff --git a/ruoyi-common/ruoyi-common-push/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/ruoyi-common/ruoyi-common-push/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports index 2e02f912d..db1bfd7a5 100644 --- a/ruoyi-common/ruoyi-common-push/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports +++ b/ruoyi-common/ruoyi-common-push/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports @@ -1,3 +1,4 @@ org.dromara.common.push.config.MessageAutoConfiguration org.dromara.common.push.config.MessageSseConfiguration org.dromara.common.push.config.MessageWebSocketConfiguration +org.dromara.common.push.config.MessageMqttConfiguration