diff --git a/yudao-module-farm/pom.xml b/yudao-module-farm/pom.xml
index f0136113..2dfa7f88 100644
--- a/yudao-module-farm/pom.xml
+++ b/yudao-module-farm/pom.xml
@@ -80,6 +80,13 @@
yudao-spring-boot-starter-mq
+
+
+ org.eclipse.paho
+ org.eclipse.paho.client.mqttv3
+ 1.2.5
+
+
cn.iocoder.boot
@@ -134,14 +141,10 @@
spring-websocket
-
-
-
-
-
-
-
-
+
+ org.springframework.boot
+ spring-boot-starter-data-mongodb
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/FarmMongoConfiguration.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/FarmMongoConfiguration.java
new file mode 100644
index 00000000..b57e0dc7
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/FarmMongoConfiguration.java
@@ -0,0 +1,15 @@
+package cn.iocoder.yudao.module.farm.config;
+
+import cn.iocoder.yudao.module.farm.dal.mongo.telemetry.DeviceTelemetryRepository;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.data.mongodb.repository.config.EnableMongoRepositories;
+
+/**
+ * Enable Mongo repositories for the farm module.
+ */
+@Configuration
+@EnableMongoRepositories(basePackageClasses = DeviceTelemetryRepository.class)
+public class FarmMongoConfiguration {
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttAutoConfiguration.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttAutoConfiguration.java
new file mode 100644
index 00000000..9195c575
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttAutoConfiguration.java
@@ -0,0 +1,48 @@
+package cn.iocoder.yudao.module.farm.config;
+
+import org.eclipse.paho.client.mqttv3.MqttClient;
+import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
+import cn.iocoder.yudao.module.farm.service.mqtt.MqttClientManager;
+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;
+
+@Configuration
+@EnableConfigurationProperties(MqttProperties.class)
+@ConditionalOnProperty(prefix = "mqtt", name = "broker")
+public class MqttAutoConfiguration {
+
+ @Bean
+ public MqttClient mqttClient(MqttProperties props) throws MqttException {
+ MqttConnectOptions options = new MqttConnectOptions();
+ options.setCleanSession(props.isCleanSession());
+ options.setKeepAliveInterval(props.getKeepAlive());
+ if (props.getUsername() != null && !props.getUsername().isEmpty()) {
+ options.setUserName(props.getUsername());
+ options.setPassword(props.getPassword() != null ? props.getPassword().toCharArray() : new char[0]);
+ }
+ MqttClient client = new MqttClient(props.getBroker(), props.getClientId(), new MemoryPersistence());
+ client.connect(options);
+ return client;
+ }
+
+ @Bean
+ @ConditionalOnProperty(prefix = "mqtt", name = "connections")
+ public java.util.Map multipleMqttClients(MqttProperties props, MqttClientManager manager) throws MqttException {
+ java.util.Map map = new java.util.HashMap<>();
+ if (props.getConnections() != null) {
+ for (MqttConnectionProperties c : props.getConnections()) {
+ MqttClient cli = manager.createAndConnect(
+ c.getName(), c.getBroker(), c.getClientId(), c.getUsername(), c.getPassword(), c.isCleanSession(), c.getKeepAlive()
+ );
+ map.put(c.getName(), cli);
+ }
+ }
+ return map;
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttConnectionProperties.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttConnectionProperties.java
new file mode 100644
index 00000000..528c53f0
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttConnectionProperties.java
@@ -0,0 +1,97 @@
+package cn.iocoder.yudao.module.farm.config;
+
+public class MqttConnectionProperties {
+
+ private String name;
+ private String broker;
+ private String deviceId;
+ private String clientId;
+ private String username;
+ private String password;
+ private boolean cleanSession = true;
+ private int keepAlive = 30;
+ private int qos = 1;
+ private boolean reconnect = true;
+
+ public String getName() {
+ return name;
+ }
+
+ public void setName(String name) {
+ this.name = name;
+ }
+
+ public String getBroker() {
+ return broker;
+ }
+
+ public void setBroker(String broker) {
+ this.broker = broker;
+ }
+
+ public String getDeviceId() {
+ return deviceId;
+ }
+
+ public void setDeviceId(String deviceId) {
+ this.deviceId = deviceId;
+ }
+
+ public String getClientId() {
+ return clientId;
+ }
+
+ public void setClientId(String clientId) {
+ this.clientId = clientId;
+ }
+
+ public String getUsername() {
+ return username;
+ }
+
+ public void setUsername(String username) {
+ this.username = username;
+ }
+
+ public String getPassword() {
+ return password;
+ }
+
+ public void setPassword(String password) {
+ this.password = password;
+ }
+
+ public boolean isCleanSession() {
+ return cleanSession;
+ }
+
+ public void setCleanSession(boolean cleanSession) {
+ this.cleanSession = cleanSession;
+ }
+
+ public int getKeepAlive() {
+ return keepAlive;
+ }
+
+ public void setKeepAlive(int keepAlive) {
+ this.keepAlive = keepAlive;
+ }
+
+ public int getQos() {
+ return qos;
+ }
+
+ public void setQos(int qos) {
+ this.qos = qos;
+ }
+
+ public boolean isReconnect() {
+ return reconnect;
+ }
+
+ public void setReconnect(boolean reconnect) {
+ this.reconnect = reconnect;
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttProperties.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttProperties.java
new file mode 100644
index 00000000..6f176418
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/config/MqttProperties.java
@@ -0,0 +1,92 @@
+package cn.iocoder.yudao.module.farm.config;
+
+import org.springframework.boot.context.properties.ConfigurationProperties;
+
+@ConfigurationProperties(prefix = "mqtt")
+public class MqttProperties {
+
+ private String broker;
+ private String clientId;
+ private String username;
+ private String password;
+ private boolean cleanSession = true;
+ private int keepAlive = 30;
+ private int qos = 1;
+ private boolean reconnect = true;
+
+ private java.util.List connections;
+
+ public String getBroker() {
+ return broker;
+ }
+
+ public void setBroker(String broker) {
+ this.broker = broker;
+ }
+
+ public String getClientId() {
+ return clientId;
+ }
+
+ public void setClientId(String clientId) {
+ this.clientId = clientId;
+ }
+
+ public String getUsername() {
+ return username;
+ }
+
+ public void setUsername(String username) {
+ this.username = username;
+ }
+
+ public String getPassword() {
+ return password;
+ }
+
+ public void setPassword(String password) {
+ this.password = password;
+ }
+
+ public boolean isCleanSession() {
+ return cleanSession;
+ }
+
+ public void setCleanSession(boolean cleanSession) {
+ this.cleanSession = cleanSession;
+ }
+
+ public int getKeepAlive() {
+ return keepAlive;
+ }
+
+ public void setKeepAlive(int keepAlive) {
+ this.keepAlive = keepAlive;
+ }
+
+ public int getQos() {
+ return qos;
+ }
+
+ public void setQos(int qos) {
+ this.qos = qos;
+ }
+
+ public boolean isReconnect() {
+ return reconnect;
+ }
+
+ public void setReconnect(boolean reconnect) {
+ this.reconnect = reconnect;
+ }
+
+ public java.util.List getConnections() {
+ return connections;
+ }
+
+ public void setConnections(java.util.List connections) {
+ this.connections = connections;
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/mqtt/MqttConnectionDO.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/mqtt/MqttConnectionDO.java
new file mode 100644
index 00000000..cf907fe0
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/dataobject/mqtt/MqttConnectionDO.java
@@ -0,0 +1,47 @@
+package cn.iocoder.yudao.module.farm.dal.dataobject.mqtt;
+
+import cn.iocoder.yudao.framework.mybatis.core.dataobject.BaseDO;
+import com.baomidou.mybatisplus.annotation.KeySequence;
+import com.baomidou.mybatisplus.annotation.TableId;
+import com.baomidou.mybatisplus.annotation.TableName;
+import lombok.*;
+
+@TableName("farm_mqtt_connection")
+@KeySequence("farm_mqtt_connection_seq")
+@Data
+@EqualsAndHashCode(callSuper = true)
+@ToString(callSuper = true)
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+public class MqttConnectionDO extends BaseDO {
+
+ @TableId
+ private Long id;
+
+ private String tenantId;
+
+ private String name;
+
+ private String deviceId;
+
+ private String broker;
+
+ private String clientId;
+
+ private String username;
+
+ private String password;
+
+ private Boolean cleanSession;
+
+ private Integer keepAlive;
+
+ private Integer qos;
+
+ private Boolean reconnect;
+
+ private Boolean status;
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mongo/telemetry/DeviceTelemetryDoc.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mongo/telemetry/DeviceTelemetryDoc.java
new file mode 100644
index 00000000..30b1bc47
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mongo/telemetry/DeviceTelemetryDoc.java
@@ -0,0 +1,28 @@
+package cn.iocoder.yudao.module.farm.dal.mongo.telemetry;
+
+import org.springframework.data.annotation.Id;
+import org.springframework.data.mongodb.core.mapping.Document;
+
+@Document(collection = "farm_device_telemetry")
+public class DeviceTelemetryDoc {
+
+ @Id
+ private String id;
+ private String tenantId;
+ private String deviceId;
+ private Long ts;
+ private String payloadJson;
+
+ public String getId() { return id; }
+ public void setId(String id) { this.id = id; }
+ public String getTenantId() { return tenantId; }
+ public void setTenantId(String tenantId) { this.tenantId = tenantId; }
+ public String getDeviceId() { return deviceId; }
+ public void setDeviceId(String deviceId) { this.deviceId = deviceId; }
+ public Long getTs() { return ts; }
+ public void setTs(Long ts) { this.ts = ts; }
+ public String getPayloadJson() { return payloadJson; }
+ public void setPayloadJson(String payloadJson) { this.payloadJson = payloadJson; }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mongo/telemetry/DeviceTelemetryRepository.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mongo/telemetry/DeviceTelemetryRepository.java
new file mode 100644
index 00000000..e7e3f102
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mongo/telemetry/DeviceTelemetryRepository.java
@@ -0,0 +1,11 @@
+package cn.iocoder.yudao.module.farm.dal.mongo.telemetry;
+
+import org.springframework.data.mongodb.repository.MongoRepository;
+
+import java.util.List;
+
+public interface DeviceTelemetryRepository extends MongoRepository {
+ List findTop100ByTenantIdAndDeviceIdOrderByTsDesc(String tenantId, String deviceId);
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/mqtt/MqttConnectionMapper.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/mqtt/MqttConnectionMapper.java
new file mode 100644
index 00000000..cea5038d
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/dal/mysql/mqtt/MqttConnectionMapper.java
@@ -0,0 +1,20 @@
+package cn.iocoder.yudao.module.farm.dal.mysql.mqtt;
+
+import cn.iocoder.yudao.framework.mybatis.core.mapper.BaseMapperX;
+import cn.iocoder.yudao.framework.mybatis.core.query.LambdaQueryWrapperX;
+import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
+import org.apache.ibatis.annotations.Mapper;
+
+@Mapper
+public interface MqttConnectionMapper extends BaseMapperX {
+
+ default MqttConnectionDO selectByDeviceIdAndTenant(String deviceId, String tenantId) {
+ return selectOne(new LambdaQueryWrapperX()
+ .eqIfPresent(MqttConnectionDO::getDeviceId, deviceId)
+ .eqIfPresent(MqttConnectionDO::getTenantId, tenantId)
+ .eq(MqttConnectionDO::getStatus, true)
+ );
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/DeviceMqttService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/DeviceMqttService.java
new file mode 100644
index 00000000..74d55f56
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/DeviceMqttService.java
@@ -0,0 +1,83 @@
+package cn.iocoder.yudao.module.farm.service.mqtt;
+
+import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
+import cn.iocoder.yudao.module.farm.dal.mysql.deviceinfo.DeviceInfoMapper;
+import org.eclipse.paho.client.mqttv3.IMqttMessageListener;
+import org.eclipse.paho.client.mqttv3.MqttClient;
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.springframework.stereotype.Service;
+
+import jakarta.annotation.Resource;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.nio.charset.StandardCharsets;
+import cn.iocoder.yudao.module.farm.service.telemetry.DeviceTelemetryService;
+
+@Service
+public class DeviceMqttService {
+
+ @Resource
+ private DeviceInfoMapper deviceInfoMapper;
+ @Resource
+ private MqttConnectionService mqttConnectionService;
+ @Resource
+ private DeviceTelemetryService deviceTelemetryService;
+
+ private final Set subscribed = ConcurrentHashMap.newKeySet();
+
+ // 上报主题
+ private String telemetryTopic(String deviceUniqueId) {
+ return "farm/" + deviceUniqueId + "/telemetry";
+ }
+
+ // 下发指令
+ private String commandTopic(String deviceUniqueId) {
+ return "farm/" + deviceUniqueId + "/cmd";
+ }
+
+ public void publishTelemetry(String deviceUniqueId, String jsonPayload) throws MqttException {
+ MqttClient client = mqttConnectionService.getClientByDeviceId(deviceUniqueId);
+ if (client == null) return;
+ client.publish(telemetryTopic(deviceUniqueId), jsonPayload.getBytes(StandardCharsets.UTF_8), 1, false);
+ }
+
+ public void sendCommand(String deviceUniqueId, String jsonCommand) throws MqttException {
+ MqttClient client = mqttConnectionService.getClientByDeviceId(deviceUniqueId);
+ if (client == null) return;
+ client.publish(commandTopic(deviceUniqueId), jsonCommand.getBytes(StandardCharsets.UTF_8), 1, false);
+ }
+
+ public void subscribeTelemetry(String deviceUniqueId, IMqttMessageListener listener) throws MqttException {
+ MqttClient client = mqttConnectionService.getClientByDeviceId(deviceUniqueId);
+ if (client == null) return;
+ String topic = telemetryTopic(deviceUniqueId);
+ String key = client.getClientId() + "|" + topic;
+ if (subscribed.add(key)) {
+ client.subscribe(topic, 1, listener);
+ }
+ }
+
+ public void subscribeCommands(String deviceUniqueId, IMqttMessageListener listener) throws MqttException {
+ MqttClient client = mqttConnectionService.getClientByDeviceId(deviceUniqueId);
+ if (client == null) return;
+ String topic = commandTopic(deviceUniqueId);
+ String key = client.getClientId() + "|" + topic;
+ if (subscribed.add(key)) {
+ client.subscribe(topic, 1, listener);
+ }
+ }
+
+ // 订阅并自动入库
+ public void subscribeTelemetryAndPersist(String deviceUniqueId) throws MqttException {
+ this.subscribeTelemetry(deviceUniqueId, (topic, message) -> {
+ String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
+ deviceTelemetryService.saveTelemetry(deviceUniqueId, payload);
+ });
+ }
+
+ public DeviceInfoDO getDevice(String deviceUniqueId) {
+ return deviceInfoMapper.selectOne(DeviceInfoDO::getDeviceUniqueId, deviceUniqueId);
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttClientManager.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttClientManager.java
new file mode 100644
index 00000000..62e65cd4
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttClientManager.java
@@ -0,0 +1,50 @@
+package cn.iocoder.yudao.module.farm.service.mqtt;
+
+import org.eclipse.paho.client.mqttv3.MqttClient;
+import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
+import org.springframework.stereotype.Service;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+
+@Service
+public class MqttClientManager {
+
+ private final Map clients = new ConcurrentHashMap<>();
+
+ public MqttClient createAndConnect(String name,
+ String broker,
+ String clientId,
+ String username,
+ String password,
+ boolean cleanSession,
+ int keepAlive) throws MqttException {
+ MqttClient client = new MqttClient(broker, clientId, new MemoryPersistence());
+ MqttConnectOptions options = new MqttConnectOptions();
+ options.setCleanSession(cleanSession);
+ options.setKeepAliveInterval(keepAlive);
+ if (username != null && !username.isEmpty()) {
+ options.setUserName(username);
+ options.setPassword(password != null ? password.toCharArray() : new char[0]);
+ }
+ client.connect(options);
+ clients.put(name, client);
+ return client;
+ }
+
+ public MqttClient get(String name) {
+ return clients.get(name);
+ }
+
+ public void disconnectAndRemove(String name) {
+ MqttClient c = clients.remove(name);
+ if (c != null) {
+ try { c.disconnect(); } catch (Exception ignored) {}
+ try { c.close(); } catch (Exception ignored) {}
+ }
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttClientService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttClientService.java
new file mode 100644
index 00000000..3d3aada0
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttClientService.java
@@ -0,0 +1,30 @@
+package cn.iocoder.yudao.module.farm.service.mqtt;
+
+import org.eclipse.paho.client.mqttv3.IMqttMessageListener;
+import org.eclipse.paho.client.mqttv3.MqttClient;
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.springframework.stereotype.Service;
+
+import jakarta.annotation.Resource;
+import java.nio.charset.StandardCharsets;
+
+@Service
+public class MqttClientService {
+
+ @Resource
+ private MqttConnectionService mqttConnectionService;
+
+ public void publish(String deviceId, String topic, String payload, int qos, boolean retained) throws MqttException {
+ MqttClient client = mqttConnectionService.getClientByDeviceId(deviceId);
+ if (client == null) return;
+ client.publish(topic, payload.getBytes(StandardCharsets.UTF_8), qos, retained);
+ }
+
+ public void subscribe(String deviceId, String topicFilter, int qos, IMqttMessageListener listener) throws MqttException {
+ MqttClient client = mqttConnectionService.getClientByDeviceId(deviceId);
+ if (client == null) return;
+ client.subscribe(topicFilter, qos, listener);
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionService.java
new file mode 100644
index 00000000..8b80a79b
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/mqtt/MqttConnectionService.java
@@ -0,0 +1,102 @@
+package cn.iocoder.yudao.module.farm.service.mqtt;
+
+import cn.iocoder.yudao.module.farm.dal.dataobject.mqtt.MqttConnectionDO;
+import cn.iocoder.yudao.module.farm.dal.mysql.mqtt.MqttConnectionMapper;
+import org.eclipse.paho.client.mqttv3.IMqttMessageListener;
+import org.eclipse.paho.client.mqttv3.MqttClient;
+import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
+import org.eclipse.paho.client.mqttv3.MqttException;
+import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
+import org.eclipse.paho.client.mqttv3.MqttMessage;
+import org.springframework.stereotype.Service;
+
+import jakarta.annotation.PreDestroy;
+import jakarta.annotation.Resource;
+import java.nio.charset.StandardCharsets;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+
+import static cn.iocoder.yudao.framework.security.core.util.SecurityFrameworkUtils.getLoginUser;
+
+@Service
+public class MqttConnectionService {
+
+ @Resource
+ private MqttConnectionMapper connectionMapper;
+
+ private final Map cache = new ConcurrentHashMap<>();
+
+ public MqttClient getClientByDeviceId(String deviceId) throws MqttException {
+ String tenantId = String.valueOf(getLoginUser().getTenantId());
+ return this.getClientByDeviceId(deviceId, tenantId);
+ }
+
+ public MqttClient getClientByDeviceId(String deviceId, String tenantId) throws MqttException {
+ String key = tenantId + ":" + deviceId;
+ MqttClient existing = cache.get(key);
+ if (existing != null) {
+ if (!existing.isConnected()) {
+ try { existing.reconnect(); } catch (Exception ignored) {}
+ }
+ if (existing.isConnected()) return existing;
+ }
+
+ MqttConnectionDO cfg = connectionMapper.selectByDeviceIdAndTenant(deviceId, tenantId);
+ if (cfg == null) {
+ return null;
+ }
+
+ MqttClient client = new MqttClient(cfg.getBroker(), cfg.getClientId(), new MemoryPersistence());
+ MqttConnectOptions options = new MqttConnectOptions();
+ options.setCleanSession(Boolean.TRUE.equals(cfg.getCleanSession()));
+ options.setAutomaticReconnect(Boolean.TRUE.equals(cfg.getReconnect()));
+ if (cfg.getKeepAlive() != null) options.setKeepAliveInterval(cfg.getKeepAlive());
+ if (cfg.getUsername() != null && !cfg.getUsername().isEmpty()) {
+ options.setUserName(cfg.getUsername());
+ options.setPassword(cfg.getPassword() != null ? cfg.getPassword().toCharArray() : new char[0]);
+ }
+ try {
+ String lwtTopic = "farm/" + deviceId + "/lwt";
+ MqttMessage lwt = new MqttMessage("offline".getBytes(StandardCharsets.UTF_8));
+ lwt.setQos(cfg.getQos() != null ? cfg.getQos() : 1);
+ lwt.setRetained(true);
+ options.setWill(lwtTopic, lwt.getPayload(), lwt.getQos(), lwt.isRetained());
+ } catch (Exception ignored) {}
+ client.connect(options);
+ cache.put(key, client);
+ return client;
+ }
+
+ public void publish(String deviceId, String topic, String payload, int qos, boolean retained) throws MqttException {
+ MqttClient c = getClientByDeviceId(deviceId);
+ if (c == null) return;
+ c.publish(topic, payload.getBytes(StandardCharsets.UTF_8), qos, retained);
+ }
+
+ public void subscribe(String deviceId, String topicFilter, int qos, IMqttMessageListener listener) throws MqttException {
+ MqttClient c = getClientByDeviceId(deviceId);
+ if (c == null) return;
+ c.subscribe(topicFilter, qos, listener);
+ }
+
+ public void reload(String deviceId) {
+ String tenantId = String.valueOf(getLoginUser().getTenantId());
+ String key = tenantId + ":" + deviceId;
+ MqttClient old = cache.remove(key);
+ if (old != null) {
+ try { old.disconnect(); } catch (Exception ignored) {}
+ try { old.close(); } catch (Exception ignored) {}
+ }
+ }
+
+ @PreDestroy
+ public void shutdown() {
+ for (Map.Entry e : cache.entrySet()) {
+ try { e.getValue().disconnect(); } catch (Exception ignored) {}
+ try { e.getValue().close(); } catch (Exception ignored) {}
+ }
+ cache.clear();
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/telemetry/DeviceTelemetryService.java b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/telemetry/DeviceTelemetryService.java
new file mode 100644
index 00000000..fdd310cf
--- /dev/null
+++ b/yudao-module-farm/src/main/java/cn/iocoder/yudao/module/farm/service/telemetry/DeviceTelemetryService.java
@@ -0,0 +1,39 @@
+package cn.iocoder.yudao.module.farm.service.telemetry;
+
+import cn.iocoder.yudao.module.farm.dal.dataobject.deviceinfo.DeviceInfoDO;
+import cn.iocoder.yudao.module.farm.dal.mongo.telemetry.DeviceTelemetryDoc;
+import cn.iocoder.yudao.module.farm.dal.mongo.telemetry.DeviceTelemetryRepository;
+import cn.iocoder.yudao.module.farm.dal.mysql.deviceinfo.DeviceInfoMapper;
+import org.springframework.stereotype.Service;
+
+import jakarta.annotation.Resource;
+
+import static cn.iocoder.yudao.framework.security.core.util.SecurityFrameworkUtils.getLoginUser;
+
+@Service
+public class DeviceTelemetryService {
+
+ @Resource
+ private DeviceTelemetryRepository telemetryRepository;
+ @Resource
+ private DeviceInfoMapper deviceInfoMapper;
+
+ public void saveTelemetry(String deviceUniqueId, String payloadJson) {
+ String tenantId = String.valueOf(getLoginUser().getTenantId());
+
+ DeviceTelemetryDoc doc = new DeviceTelemetryDoc();
+ doc.setTenantId(tenantId);
+ doc.setDeviceId(deviceUniqueId);
+ doc.setTs(System.currentTimeMillis());
+ doc.setPayloadJson(payloadJson);
+ telemetryRepository.save(doc);
+
+ DeviceInfoDO device = deviceInfoMapper.selectOne(DeviceInfoDO::getDeviceUniqueId, deviceUniqueId);
+ if (device != null) {
+ device.setDataImage(payloadJson);
+ deviceInfoMapper.updateById(device);
+ }
+ }
+}
+
+
diff --git a/yudao-module-farm/src/main/resources/mapper/personnelmanagement/SystemUserRoleMapper.xml b/yudao-module-farm/src/main/resources/mapper/personnelmanagement/SystemUserRoleMapper.xml
new file mode 100644
index 00000000..15a4d650
--- /dev/null
+++ b/yudao-module-farm/src/main/resources/mapper/personnelmanagement/SystemUserRoleMapper.xml
@@ -0,0 +1,13 @@
+
+
+
+
+
+
+
+
\ No newline at end of file
diff --git a/yudao-server/src/main/resources/application-dev.yaml b/yudao-server/src/main/resources/application-dev.yaml
index 131d3937..4a2e2a15 100644
--- a/yudao-server/src/main/resources/application-dev.yaml
+++ b/yudao-server/src/main/resources/application-dev.yaml
@@ -122,6 +122,7 @@ lock4j:
acquire-timeout: 3000 # 获取分布式锁超时时间,默认为 3000 毫秒
expire: 30000 # 分布式锁的超时时间,默认为 30 毫秒
+
--- #################### 监控相关配置 ####################
# Actuator 监控端点的配置项