farm模块-mqtt优化
This commit is contained in:
@@ -80,6 +80,13 @@
|
||||
<artifactId>yudao-spring-boot-starter-mq</artifactId>
|
||||
</dependency>
|
||||
|
||||
<!-- MQTT client -->
|
||||
<dependency>
|
||||
<groupId>org.eclipse.paho</groupId>
|
||||
<artifactId>org.eclipse.paho.client.mqttv3</artifactId>
|
||||
<version>1.2.5</version>
|
||||
</dependency>
|
||||
|
||||
<!-- Test 测试相关 -->
|
||||
<dependency>
|
||||
<groupId>cn.iocoder.boot</groupId>
|
||||
@@ -134,14 +141,10 @@
|
||||
<artifactId>spring-websocket</artifactId>
|
||||
</dependency>
|
||||
<!-- MongoDB -->
|
||||
<!-- <dependency>-->
|
||||
<!-- <groupId>org.springframework.boot</groupId>-->
|
||||
<!-- <artifactId>spring-boot-starter-data-mongodb</artifactId>-->
|
||||
<!-- </dependency>-->
|
||||
<!-- <dependency>-->
|
||||
<!-- <groupId>org.springframework.data</groupId>-->
|
||||
<!-- <artifactId>spring-data-mongodb</artifactId>-->
|
||||
<!-- </dependency>-->
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-data-mongodb</artifactId>
|
||||
</dependency>
|
||||
|
||||
<!-- Validation for @Valid/@NotNull etc. -->
|
||||
<!-- <dependency>-->
|
||||
|
||||
@@ -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 {
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String, MqttClient> multipleMqttClients(MqttProperties props, MqttClientManager manager) throws MqttException {
|
||||
java.util.Map<String, MqttClient> 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;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<MqttConnectionProperties> 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<MqttConnectionProperties> getConnections() {
|
||||
return connections;
|
||||
}
|
||||
|
||||
public void setConnections(java.util.List<MqttConnectionProperties> connections) {
|
||||
this.connections = connections;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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; }
|
||||
}
|
||||
|
||||
|
||||
@@ -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<DeviceTelemetryDoc, String> {
|
||||
List<DeviceTelemetryDoc> findTop100ByTenantIdAndDeviceIdOrderByTsDesc(String tenantId, String deviceId);
|
||||
}
|
||||
|
||||
|
||||
@@ -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<MqttConnectionDO> {
|
||||
|
||||
default MqttConnectionDO selectByDeviceIdAndTenant(String deviceId, String tenantId) {
|
||||
return selectOne(new LambdaQueryWrapperX<MqttConnectionDO>()
|
||||
.eqIfPresent(MqttConnectionDO::getDeviceId, deviceId)
|
||||
.eqIfPresent(MqttConnectionDO::getTenantId, tenantId)
|
||||
.eq(MqttConnectionDO::getStatus, true)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String> 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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String, MqttClient> 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) {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String, MqttClient> 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<String, MqttClient> e : cache.entrySet()) {
|
||||
try { e.getValue().disconnect(); } catch (Exception ignored) {}
|
||||
try { e.getValue().close(); } catch (Exception ignored) {}
|
||||
}
|
||||
cache.clear();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
|
||||
<mapper namespace="cn.iocoder.yudao.module.farm.dal.mysql.personnelmanagement.SystemUserRoleMapper">
|
||||
|
||||
<!--
|
||||
一般情况下,尽可能使用 Mapper 进行 CRUD 增删改查即可。
|
||||
无法满足的场景,例如说多表关联查询,才使用 XML 编写 SQL。
|
||||
代码生成器暂时只生成 Mapper XML 文件本身,更多推荐 MybatisX 快速开发插件来生成查询。
|
||||
文档可见:https://www.iocoder.cn/MyBatis/x-plugins/
|
||||
-->
|
||||
|
||||
|
||||
</mapper>
|
||||
@@ -122,6 +122,7 @@ lock4j:
|
||||
acquire-timeout: 3000 # 获取分布式锁超时时间,默认为 3000 毫秒
|
||||
expire: 30000 # 分布式锁的超时时间,默认为 30 毫秒
|
||||
|
||||
|
||||
--- #################### 监控相关配置 ####################
|
||||
|
||||
# Actuator 监控端点的配置项
|
||||
|
||||
Reference in New Issue
Block a user