fix:优化物联设备mqtt订阅

main
HuangHuiKang 4 days ago
parent 0e5d780de8
commit 566acc12eb

@ -87,7 +87,7 @@ public class DefaultEmqConfig {
* @return
* @throws Exception
*/
@Bean
@Bean(destroyMethod = "")
public MqttClient mqttClient(MqttConnectOptions options, DefaultEmqProperties emqProperties, DefaultBizTopicSet defaultBizTopicSet, ApplicationContext applicationContext) throws Exception {
MqttClient mqttClient = new MqttClient(emqProperties.getBroker(), emqProperties.getClientId(), new MemoryPersistence());
mqttClient.setCallback(new MqttCallbackImpl(defaultBizTopicSet.getTopicMap(), mqttClient, options));

@ -0,0 +1,30 @@
package cn.iocoder.yudao.module.iot.framework.mqtt.config;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.MqttClient;
import org.springframework.stereotype.Component;
import javax.annotation.PreDestroy;
import javax.annotation.Resource;
@Slf4j
@Component
public class MqttClientShutdown {
@Resource
private MqttClient mqttClient;
@PreDestroy
public void shutdown() {
try {
if (mqttClient.isConnected()) {
mqttClient.disconnectForcibly(1000, 1000);
}
mqttClient.close();
log.info("MQTT客户端已关闭");
} catch (Exception e) {
log.warn("MQTT客户端关闭失败: {}", e.getMessage());
}
}
}

@ -0,0 +1,164 @@
package cn.iocoder.yudao.module.iot.framework.mqtt.consumer;
import cn.iocoder.yudao.module.iot.controller.admin.devicemodelrules.vo.PointRulesRespVO;
import cn.iocoder.yudao.module.iot.dal.dataobject.device.DeviceDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.devicecontactmodel.DeviceContactModelDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.devicepointrules.DevicePointRulesDO;
import cn.iocoder.yudao.module.iot.dal.mysql.device.DeviceMapper;
import cn.iocoder.yudao.module.iot.dal.mysql.devicecontactmodel.DeviceContactModelMapper;
import cn.iocoder.yudao.module.iot.dal.mysql.devicepointrules.DevicePointRulesMapper;
import com.alibaba.fastjson.JSON;
import com.baomidou.mybatisplus.core.toolkit.CollectionUtils;
import com.baomidou.mybatisplus.core.toolkit.Wrappers;
import lombok.AllArgsConstructor;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.stream.Collectors;
@Slf4j
@Component
public class IotMqttRuntimeCache {
@Resource
@Lazy
private DeviceMapper deviceMapper;
@Resource
private DeviceContactModelMapper deviceContactModelMapper;
@Resource
private DevicePointRulesMapper devicePointRulesMapper;
private final ConcurrentMap<String, Optional<DeviceDO>> deviceByTopic = new ConcurrentHashMap<>();
private final ConcurrentMap<Long, List<DeviceContactModelDO>> pointsByDeviceId = new ConcurrentHashMap<>();
private final ConcurrentMap<Long, Map<String, List<CachedPointRule>>> pointRulesByDeviceId = new ConcurrentHashMap<>();
private final ConcurrentMap<Long, Optional<DevicePointRulesDO>> countRuleByDeviceId = new ConcurrentHashMap<>();
public DeviceDO getDeviceByTopic(String topic) {
if (StringUtils.isBlank(topic)) {
return null;
}
return deviceByTopic.computeIfAbsent(topic, this::loadDeviceByTopic).orElse(null);
}
public List<DeviceContactModelDO> getDevicePoints(Long deviceId) {
if (deviceId == null) {
return Collections.emptyList();
}
return pointsByDeviceId.computeIfAbsent(deviceId, this::loadDevicePoints);
}
public List<CachedPointRule> getPointRules(Long deviceId, String attributeCode) {
if (deviceId == null || StringUtils.isBlank(attributeCode)) {
return Collections.emptyList();
}
return pointRulesByDeviceId.computeIfAbsent(deviceId, this::loadPointRulesByAttributeCode)
.getOrDefault(attributeCode, Collections.emptyList());
}
public DevicePointRulesDO getLatestCountRule(Long deviceId) {
if (deviceId == null) {
return null;
}
return countRuleByDeviceId.computeIfAbsent(deviceId, this::loadLatestCountRule).orElse(null);
}
public void refreshDevice(Long deviceId) {
if (deviceId == null) {
return;
}
pointsByDeviceId.remove(deviceId);
pointRulesByDeviceId.remove(deviceId);
countRuleByDeviceId.remove(deviceId);
deviceByTopic.entrySet().removeIf(entry -> entry.getValue().map(device -> deviceId.equals(device.getId())).orElse(false));
}
public void refreshTopic(String topic) {
if (StringUtils.isNotBlank(topic)) {
deviceByTopic.remove(topic);
}
}
public void clearAll() {
deviceByTopic.clear();
pointsByDeviceId.clear();
pointRulesByDeviceId.clear();
countRuleByDeviceId.clear();
}
private Optional<DeviceDO> loadDeviceByTopic(String topic) {
return Optional.ofNullable(deviceMapper.selectOne(Wrappers.<DeviceDO>lambdaQuery()
.eq(DeviceDO::getTopic, topic)
.last("LIMIT 1")));
}
private List<DeviceContactModelDO> loadDevicePoints(Long deviceId) {
List<DeviceContactModelDO> points = deviceContactModelMapper.selectList(Wrappers.<DeviceContactModelDO>lambdaQuery()
.eq(DeviceContactModelDO::getDeviceId, deviceId));
if (CollectionUtils.isEmpty(points)) {
return Collections.emptyList();
}
return Collections.unmodifiableList(new ArrayList<>(points));
}
private Map<String, List<CachedPointRule>> loadPointRulesByAttributeCode(Long deviceId) {
List<DevicePointRulesDO> rules = devicePointRulesMapper.selectList(Wrappers.<DevicePointRulesDO>lambdaQuery()
.eq(DevicePointRulesDO::getDeviceId, deviceId)
.orderByDesc(DevicePointRulesDO::getCreateTime));
if (CollectionUtils.isEmpty(rules)) {
return Collections.emptyMap();
}
List<CachedPointRule> parsedRules = new ArrayList<>();
for (DevicePointRulesDO rule : rules) {
if (StringUtils.isBlank(rule.getFieldRule())) {
continue;
}
try {
List<PointRulesRespVO> pointRules = JSON.parseArray(rule.getFieldRule(), PointRulesRespVO.class);
if (CollectionUtils.isEmpty(pointRules)) {
continue;
}
for (PointRulesRespVO pointRule : pointRules) {
if (StringUtils.isBlank(pointRule.getCode())) {
continue;
}
parsedRules.add(new CachedPointRule(rule, pointRule));
}
} catch (Exception e) {
log.warn("MQTT规则解析失败 deviceId={}, ruleId={}, reason={}", deviceId, rule.getId(), e.getMessage());
}
}
if (parsedRules.isEmpty()) {
return Collections.emptyMap();
}
return parsedRules.stream().collect(Collectors.groupingBy(rule -> rule.getPointRule().getCode()));
}
private Optional<DevicePointRulesDO> loadLatestCountRule(Long deviceId) {
return Optional.ofNullable(devicePointRulesMapper.selectOne(Wrappers.<DevicePointRulesDO>lambdaQuery()
.eq(DevicePointRulesDO::getDeviceId, deviceId)
.eq(DevicePointRulesDO::getIdentifier, "COUNT")
.orderByDesc(DevicePointRulesDO::getUpdateTime)
.last("LIMIT 1")));
}
@Getter
@AllArgsConstructor
public static class CachedPointRule {
private final DevicePointRulesDO rule;
private final PointRulesRespVO pointRule;
}
}

@ -5,6 +5,7 @@ import cn.iocoder.yudao.framework.tenant.core.context.TenantContextHolder;
import cn.iocoder.yudao.module.iot.controller.admin.device.enums.DeviceBasicStatusEnum;
import cn.iocoder.yudao.module.iot.controller.admin.device.enums.DeviceStatusEnum;
import cn.iocoder.yudao.module.iot.controller.admin.devicemodelrules.vo.PointRulesRespVO;
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
import cn.iocoder.yudao.module.iot.dal.dataobject.device.DeviceDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.devicecontactmodel.DeviceContactModelDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.deviceoperationrecord.DeviceOperationRecordDO;
@ -12,14 +13,12 @@ import cn.iocoder.yudao.module.iot.dal.dataobject.devicepointrules.DevicePointRu
import cn.iocoder.yudao.module.iot.dal.dataobject.devicewarinningrecord.DeviceWarinningRecordDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.iotorganization.IotOrganizationDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.mqttrecord.MqttRecordDO;
import cn.iocoder.yudao.module.iot.dal.mysql.device.DeviceMapper;
import cn.iocoder.yudao.module.iot.dal.mysql.devicecontactmodel.DeviceContactModelMapper;
import cn.iocoder.yudao.module.iot.dal.mysql.deviceoperationrecord.DeviceOperationRecordMapper;
import cn.iocoder.yudao.module.iot.dal.mysql.devicepointrules.DevicePointRulesMapper;
import cn.iocoder.yudao.module.iot.dal.mysql.devicewarinningrecord.DeviceWarinningRecordMapper;
import cn.iocoder.yudao.module.iot.dal.mysql.mqttrecord.MqttRecordMapper;
import cn.iocoder.yudao.module.iot.framework.constant.Constants;
import cn.iocoder.yudao.module.iot.framework.mqtt.common.SuperConsumer;
import cn.iocoder.yudao.module.iot.framework.mqtt.consumer.IotMqttRuntimeCache.CachedPointRule;
import cn.iocoder.yudao.module.iot.framework.mqtt.consumer.impl.AsyncService;
import cn.iocoder.yudao.module.iot.framework.mqtt.entity.MqttData;
import cn.iocoder.yudao.module.iot.framework.mqtt.utils.DateUtils;
@ -32,16 +31,13 @@ import cn.iocoder.yudao.module.iot.service.mqttrecord.MqttRecordService;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.toolkit.CollectionUtils;
import com.baomidou.mybatisplus.core.toolkit.Wrappers;
import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@ -58,9 +54,6 @@ public class MqttDataHandler extends SuperConsumer<String> {
@Resource
private IotOrganizationService organizationService;
@Resource
@Lazy
private DeviceMapper deviceMapper;
@Resource
private AsyncService asyncService;
@Resource
@ -71,17 +64,14 @@ public class MqttDataHandler extends SuperConsumer<String> {
@Resource
private DeviceOperationRecordMapper deviceOperationRecordMapper;
@Resource
private DeviceContactModelMapper deviceContactModelMapper;
@Resource
private TDengineService tDengineService;
@Resource
private DevicePointRulesMapper devicePointRulesMapper;
private DeviceWarinningRecordMapper deviceWarinningRecordMapper;
@Resource
private DeviceWarinningRecordMapper deviceWarinningRecordMapper;
private IotMqttRuntimeCache mqttRuntimeCache;
private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
@ -187,10 +177,10 @@ public class MqttDataHandler extends SuperConsumer<String> {
Map.class
);
DeviceDO deviceDO = deviceMapper.selectOne(Wrappers.<DeviceDO>lambdaQuery().eq(DeviceDO::getTopic,topic));
log.info("getDeviceByMqttTopic参数{}", topic);
DeviceDO deviceDO = mqttRuntimeCache.getDeviceByTopic(topic);
log.debug("getDeviceByMqttTopic参数{}", topic);
if (deviceDO == null) {
log.info("getDeviceByMqttTopic查询出来deviceDO为空");
log.debug("getDeviceByMqttTopic查询出来deviceDO为空");
return;
}
@ -214,7 +204,7 @@ public class MqttDataHandler extends SuperConsumer<String> {
Long deviceId = device.getId();
// 1. 查询点位配置
List<DeviceContactModelDO> points = getDevicePoints(deviceId);
List<DeviceContactModelDO> points = mqttRuntimeCache.getDevicePoints(deviceId);
if (CollectionUtils.isEmpty(points)) {
@ -258,7 +248,7 @@ public class MqttDataHandler extends SuperConsumer<String> {
log.info("设备 {} MQTT 数据点位数量 {}", deviceId, varListMap.size());
log.debug("设备 {} MQTT 数据点位数量 {}", deviceId, varListMap.size());
int successCount = 0;
@ -272,9 +262,10 @@ public class MqttDataHandler extends SuperConsumer<String> {
String code = point.getAttributeCode();
Object value = varListMap.get(code);
DeviceContactModelDO pointData = BeanUtils.toBean(point, DeviceContactModelDO.class);
if (value == null) {
validDataList.add(point);
validDataList.add(pointData);
continue;
}
@ -285,14 +276,14 @@ public class MqttDataHandler extends SuperConsumer<String> {
processedValue,
code,
device,
point.getId()
pointData
);
point.setAddressValue(processedValue);
pointData.setAddressValue(processedValue);
successCount++;
validDataList.add(point);
validDataList.add(pointData);
} catch (Exception e) {
@ -320,13 +311,7 @@ public class MqttDataHandler extends SuperConsumer<String> {
// handleCapacityFormula
private void handleCapacityFormula(DeviceDO device, Map<String, Object> varListMap) {
DevicePointRulesDO formulaRule = devicePointRulesMapper.selectOne(
Wrappers.<DevicePointRulesDO>lambdaQuery()
.eq(DevicePointRulesDO::getDeviceId, device.getId())
.eq(DevicePointRulesDO::getIdentifier, "COUNT")
.orderByDesc(DevicePointRulesDO::getUpdateTime)
.last("limit 1")
);
DevicePointRulesDO formulaRule = mqttRuntimeCache.getLatestCountRule(device.getId());
if (formulaRule == null || StringUtils.isBlank(formulaRule.getFieldRule())) {
return;
}
@ -354,46 +339,6 @@ public class MqttDataHandler extends SuperConsumer<String> {
private DevicePointRulesDO getDevicePointRules(Long deviceId) {
List<DevicePointRulesDO> list =
devicePointRulesMapper.selectList(
Wrappers.<DevicePointRulesDO>lambdaQuery()
.eq(DevicePointRulesDO::getDeviceId, deviceId)
.eq(DevicePointRulesDO::getIdentifier, "RUNNING")
.orderByDesc(DevicePointRulesDO::getCreateTime)
.last("LIMIT 1")
);
if (CollectionUtils.isEmpty(list)) {
log.info("设备 {} 未找到 RUNNING 规则", deviceId);
return null;
}
DevicePointRulesDO rule = list.get(0);
log.info("设备 {} 使用 RUNNING 规则规则ID={}, 创建时间={}",
deviceId,
rule.getId(),
rule.getCreateTime());
return rule;
}
/**
*
*/
private List<DeviceContactModelDO> getDevicePoints(Long deviceId) {
LambdaQueryWrapper<DeviceContactModelDO> query = new LambdaQueryWrapper<>();
query.eq(DeviceContactModelDO::getDeviceId, deviceId);
return deviceContactModelMapper.selectList(query);
}
/**
* OPC
*/
@ -423,7 +368,7 @@ public class MqttDataHandler extends SuperConsumer<String> {
boolean isSuccess = tDengineService.newInsertDeviceData(deviceId, dataList);
if (isSuccess) {
log.info("设备 {} 数据入库成功,总数: {},有效: {}",
log.debug("设备 {} 数据入库成功,总数: {},有效: {}",
deviceId, dataList.size(), successCount);
} else {
log.error("设备 {} 数据入库失败", deviceId);
@ -433,60 +378,31 @@ public class MqttDataHandler extends SuperConsumer<String> {
}
}
private void judgmentRules(String processedValue, String attributeCode, DeviceDO device, Long modelId) {
private void judgmentRules(String processedValue, String attributeCode, DeviceDO device, DeviceContactModelDO point) {
if (StringUtils.isBlank(processedValue)) {
log.warn("待判断的值为空编码attributeCode: {}, deviceId: {}", attributeCode, device.getId());
// return;
}
// 1. 查询设备规则
List<DevicePointRulesDO> devicePointRulesDOList = devicePointRulesMapper.selectList(
Wrappers.<DevicePointRulesDO>lambdaQuery()
.eq(DevicePointRulesDO::getDeviceId, device.getId()).orderByDesc(DevicePointRulesDO::getCreateTime));
if (CollectionUtils.isEmpty(devicePointRulesDOList)) {
List<CachedPointRule> cachedRules = mqttRuntimeCache.getPointRules(device.getId(), attributeCode);
if (CollectionUtils.isEmpty(cachedRules)) {
log.debug("设备 {} 未配置规则", device.getId());
return;
}
// 2. 遍历规则
for (DevicePointRulesDO devicePointRulesDO : devicePointRulesDOList) {
if (StringUtils.isBlank(devicePointRulesDO.getFieldRule())) {
continue;
}
// 3. 解析规则列表
List<PointRulesRespVO> pointRulesVOList = JSON.parseArray(
devicePointRulesDO.getFieldRule(), PointRulesRespVO.class);
if (CollectionUtils.isEmpty(pointRulesVOList)) {
continue;
}
// 4. 找到对应modelId的规则并进行判断
for (PointRulesRespVO pointRulesRespVO : pointRulesVOList) {
if (pointRulesRespVO.getCode() != null &&
pointRulesRespVO.getCode().equals(attributeCode)) {
boolean matched = matchRule(processedValue, pointRulesRespVO);
if (matched) {
log.info("规则匹配成功: modelId={}, value={}, rule={}",
attributeCode, processedValue,
JSON.toJSONString(pointRulesRespVO));
// 执行匹配成功后的逻辑
handleMatchedSuccessRule(devicePointRulesDO, pointRulesRespVO, processedValue, device, attributeCode, modelId);
break;
} else {
log.debug("规则不匹配: modelId={}, value={}, rule={}",
attributeCode, processedValue,
JSON.toJSONString(pointRulesRespVO));
// 执行匹配失败后的逻辑
// handleMatchedFailureRule(devicePointRulesDO, pointRulesRespVO, processedValue, device, attributeCode);
}
}
for (CachedPointRule cachedRule : cachedRules) {
PointRulesRespVO pointRulesRespVO = cachedRule.getPointRule();
boolean matched = matchRule(processedValue, pointRulesRespVO);
if (matched) {
log.debug("规则匹配成功: modelId={}, value={}, rule={}",
attributeCode, processedValue,
JSON.toJSONString(pointRulesRespVO));
handleMatchedSuccessRule(cachedRule.getRule(), pointRulesRespVO, processedValue, device, attributeCode, point);
break;
} else {
log.debug("规则不匹配: modelId={}, value={}, rule={}",
attributeCode, processedValue,
JSON.toJSONString(pointRulesRespVO));
}
}
}
@ -512,9 +428,8 @@ public class MqttDataHandler extends SuperConsumer<String> {
String processedValue,
DeviceDO device,
String attributeCode,
Long modelId) {
DeviceContactModelDO deviceContactModelDO = deviceContactModelMapper.selectById(modelId);
if (deviceContactModelDO == null) {
DeviceContactModelDO point) {
if (point == null) {
return;
}
@ -541,13 +456,13 @@ public class MqttDataHandler extends SuperConsumer<String> {
DeviceWarinningRecordDO deviceWarinningRecordDO = new DeviceWarinningRecordDO();
deviceWarinningRecordDO.setDeviceId(device.getId());
deviceWarinningRecordDO.setModelId(modelId);
deviceWarinningRecordDO.setModelId(point.getId());
deviceWarinningRecordDO.setRule(pointRulesRespVO.getRule());
deviceWarinningRecordDO.setAlarmLevel(devicePointRulesDO.getAlarmLevel());
deviceWarinningRecordDO.setAddressValue(processedValue);
deviceWarinningRecordDO.setRuleId(devicePointRulesDO.getId());
deviceWarinningRecordDO.setDeviceName(device.getDeviceName());
deviceWarinningRecordDO.setModelName(deviceContactModelDO.getAttributeName());
deviceWarinningRecordDO.setModelName(point.getAttributeName());
deviceWarinningRecordDO.setRuleName(devicePointRulesDO.getFieldName());
//TODO 创建人和更新人为内置默认管理员
deviceWarinningRecordDO.setCreator("1");

@ -1,7 +1,9 @@
package cn.iocoder.yudao.module.iot.framework.mqtt.utils;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
/**
* @author jie
@ -10,5 +12,12 @@ public class ThreadUtils {
/**
* 线
*/
public static ExecutorService executorService = Executors.newFixedThreadPool(50);
public static ExecutorService executorService = new ThreadPoolExecutor(
4,
8,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(2000),
new ThreadPoolExecutor.CallerRunsPolicy()
);
}

@ -49,6 +49,7 @@ import cn.iocoder.yudao.module.iot.dal.mysql.device.DeviceAttributeMapper;
import cn.iocoder.yudao.module.iot.framework.mqtt.config.DefaultBizTopicSet;
import cn.iocoder.yudao.module.iot.framework.mqtt.consumer.IMqttservice;
import cn.iocoder.yudao.module.iot.framework.mqtt.consumer.IotMqttRuntimeCache;
import cn.iocoder.yudao.module.iot.util.MapListStatsCalculator;
import com.alibaba.fastjson.JSON;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
@ -70,6 +71,8 @@ import org.springframework.transaction.annotation.Transactional;
import org.springframework.validation.annotation.Validated;
import javax.annotation.Resource;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.sql.Timestamp;
import java.text.SimpleDateFormat;
import java.time.LocalDate;
@ -96,6 +99,8 @@ import static cn.iocoder.yudao.module.system.enums.ErrorCodeConstants.IMPORT_LIS
@Slf4j
public class DeviceServiceImpl implements DeviceService {
private static final int ADDRESS_VALUE_SCALE = 6;
@Resource
private DeviceMapper deviceMapper;
@Resource
@ -143,6 +148,9 @@ public class DeviceServiceImpl implements DeviceService {
@Resource
private DefaultBizTopicSet defaultBizTopicSet;
@Resource
private IotMqttRuntimeCache mqttRuntimeCache;
@Resource
private ErpProductUnitMapper productUnitMapper;
@ -195,6 +203,8 @@ public class DeviceServiceImpl implements DeviceService {
//新增点位规则模板
insertTemplatePoint(createReqVO, device);
mqttRuntimeCache.refreshDevice(device.getId());
mqttRuntimeCache.refreshTopic(device.getTopic());
return device;
}
@ -285,8 +295,14 @@ public class DeviceServiceImpl implements DeviceService {
validTopicExists(updateReqVO.getTopic(), updateReqVO.getId());
// 更新
DeviceDO oldDevice = deviceMapper.selectById(updateReqVO.getId());
DeviceDO updateObj = BeanUtils.toBean(updateReqVO, DeviceDO.class);
deviceMapper.updateById(updateObj);
mqttRuntimeCache.refreshDevice(updateReqVO.getId());
if (oldDevice != null) {
mqttRuntimeCache.refreshTopic(oldDevice.getTopic());
}
mqttRuntimeCache.refreshTopic(updateReqVO.getTopic());
}
private void validTopicExists(String topic, Long currentDeviceId) {
@ -649,8 +665,12 @@ public class DeviceServiceImpl implements DeviceService {
if (value == null) return null;
try {
Double result = Double.parseDouble(value.toString());
return ratio != null ? result * ratio : result;
BigDecimal result = new BigDecimal(value.toString().trim());
if (ratio != null) {
result = result.multiply(BigDecimal.valueOf(ratio));
}
result = result.setScale(ADDRESS_VALUE_SCALE, RoundingMode.HALF_UP).stripTrailingZeros();
return result.scale() < 0 ? result.setScale(0) : result;
} catch (NumberFormatException e) {
return value; // 转换失败返回原值
}
@ -2167,6 +2187,8 @@ public class DeviceServiceImpl implements DeviceService {
// 1. 更新设备启用状态
deviceDO.setIsEnable(enabled);
deviceMapper.updateById(deviceDO);
mqttRuntimeCache.refreshDevice(deviceDO.getId());
mqttRuntimeCache.refreshTopic(deviceDO.getTopic());
String topic = deviceDO.getTopic();

@ -37,6 +37,8 @@ import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.stream.Collectors;
import static cn.iocoder.yudao.framework.common.exception.util.ServiceExceptionUtil.exception;
@ -53,6 +55,8 @@ public class TDengineService {
@Value("${spring.datasource.dynamic.datasource.tdengine.url:}")
private String tdengineUrl;
private final ConcurrentMap<Long, Set<String>> tdColumnCache = new ConcurrentHashMap<>();
private String getTdengineDbName() {
if (StrUtil.isBlank(tdengineUrl)) {
return DEFAULT_TDENGINE_DATABASE;
@ -880,6 +884,7 @@ public class TDengineService {
try {
jdbcTemplate.execute(alterSql);
markTDColumnExists(deviceId, attributeCode);
log.info("TDengine 表新增列成功: table={}, column={}, type={}",
tableName, attributeCode, tdType);
@ -903,7 +908,7 @@ public class TDengineService {
}
String tableName = tdengineTable("d_" + deviceId);
Set<String> existingColumns = getTDColumnNames(tableName);
Set<String> existingColumns = new HashSet<>(loadTDColumns(deviceId));
int failureCount = 0;
for (DeviceContactModelDO deviceContactModel : deviceContactModels) {
@ -919,9 +924,11 @@ public class TDengineService {
try {
AddTDDatabaseColumn(deviceId, deviceContactModel);
existingColumns.add(attributeCode.toLowerCase());
markTDColumnExists(deviceId, attributeCode);
} catch (Exception e) {
if (columnExists(deviceId, attributeCode)) {
existingColumns.add(attributeCode.toLowerCase());
markTDColumnExists(deviceId, attributeCode);
continue;
}
failureCount++;
@ -947,6 +954,33 @@ public class TDengineService {
}
}
private Set<String> loadTDColumns(Long deviceId) {
return tdColumnCache.computeIfAbsent(deviceId, key -> ConcurrentHashMap.newKeySet()).isEmpty()
? reloadTDColumns(deviceId)
: tdColumnCache.get(deviceId);
}
private Set<String> reloadTDColumns(Long deviceId) {
Set<String> columns = ConcurrentHashMap.newKeySet();
columns.addAll(getTDColumnNames(tdengineTable("d_" + deviceId)));
tdColumnCache.put(deviceId, columns);
return columns;
}
private void markTDColumnExists(Long deviceId, String columnName) {
if (deviceId == null || StrUtil.isBlank(columnName)) {
return;
}
tdColumnCache.computeIfAbsent(deviceId, key -> ConcurrentHashMap.newKeySet())
.add(columnName.toLowerCase());
}
public void refreshTDColumnCache(Long deviceId) {
if (deviceId != null) {
tdColumnCache.remove(deviceId);
}
}
/**
@ -1289,12 +1323,24 @@ public class TDengineService {
return false;
}
String normalizedColumnName = columnName.toLowerCase();
Set<String> cachedColumns = loadTDColumns(deviceId);
if (cachedColumns.contains(normalizedColumnName)) {
return true;
}
Set<String> reloadedColumns = reloadTDColumns(deviceId);
if (reloadedColumns.contains(normalizedColumnName)) {
return true;
}
String tableName = tdengineTable("d_" + deviceId);
try {
// 方法1直接尝试查询该列
String testSql = "SELECT " + columnName + " FROM " + tableName + " LIMIT 0";
jdbcTemplate.execute(testSql);
markTDColumnExists(deviceId, columnName);
return true; // 执行成功,说明列存在
} catch (Exception e) {

@ -4,6 +4,7 @@ import cn.hutool.core.util.StrUtil;
import cn.iocoder.yudao.module.iot.dal.dataobject.devicecontactmodel.DeviceContactModelDO;
import cn.iocoder.yudao.module.iot.dal.dataobject.devicemodelattribute.DeviceModelAttributeDO;
import cn.iocoder.yudao.module.iot.dal.mysql.deviceattributetype.DeviceAttributeTypeMapper;
import cn.iocoder.yudao.module.iot.framework.mqtt.consumer.IotMqttRuntimeCache;
import cn.iocoder.yudao.module.iot.service.device.TDengineService;
import com.baomidou.mybatisplus.core.toolkit.Wrappers;
import org.springframework.stereotype.Service;
@ -43,6 +44,8 @@ public class DeviceContactModelServiceImpl implements DeviceContactModelService
private DeviceAttributeTypeMapper deviceAttributeTypeMapper;
@Resource
private TDengineService tDengineService;
@Resource
private IotMqttRuntimeCache mqttRuntimeCache;
@Override
public Long createDeviceContactModel(DeviceContactModelSaveReqVO createReqVO) {
@ -70,6 +73,7 @@ public class DeviceContactModelServiceImpl implements DeviceContactModelService
//新增td数据库列
tDengineService.AddTDDatabaseColumn(createReqVO.getDeviceId(),deviceContactModel);
refreshMqttRuntimeCache(createReqVO.getDeviceId());
// 返回
return deviceContactModel.getId();
@ -79,9 +83,12 @@ public class DeviceContactModelServiceImpl implements DeviceContactModelService
public void updateDeviceContactModel(DeviceContactModelSaveReqVO updateReqVO) {
// 校验存在
validateDeviceContactModelExists(updateReqVO.getId());
DeviceContactModelDO oldPoint = deviceContactModelMapper.selectById(updateReqVO.getId());
// 更新
DeviceContactModelDO updateObj = BeanUtils.toBean(updateReqVO, DeviceContactModelDO.class);
deviceContactModelMapper.updateById(updateObj);
refreshMqttRuntimeCache(oldPoint == null ? null : oldPoint.getDeviceId());
refreshMqttRuntimeCache(updateReqVO.getDeviceId());
}
@Override
@ -108,7 +115,16 @@ public class DeviceContactModelServiceImpl implements DeviceContactModelService
// 删除
deviceContactModelMapper.deleteById(id);
refreshMqttRuntimeCache(deviceId);
}
}
private void refreshMqttRuntimeCache(Long deviceId) {
if (deviceId == null) {
return;
}
mqttRuntimeCache.refreshDevice(deviceId);
tDengineService.refreshTDColumnCache(deviceId);
}
/**

@ -10,6 +10,7 @@ import org.springframework.transaction.annotation.Transactional;
import java.util.*;
import cn.iocoder.yudao.module.iot.controller.admin.devicepointrules.vo.*;
import cn.iocoder.yudao.module.iot.dal.dataobject.devicepointrules.DevicePointRulesDO;
import cn.iocoder.yudao.module.iot.framework.mqtt.consumer.IotMqttRuntimeCache;
import cn.iocoder.yudao.framework.common.pojo.PageResult;
import cn.iocoder.yudao.framework.common.pojo.PageParam;
import cn.iocoder.yudao.framework.common.util.object.BeanUtils;
@ -30,12 +31,15 @@ public class DevicePointRulesServiceImpl implements DevicePointRulesService {
@Resource
private DevicePointRulesMapper devicePointRulesMapper;
@Resource
private IotMqttRuntimeCache mqttRuntimeCache;
@Override
public Long createDevicePointRules(DevicePointRulesSaveReqVO createReqVO) {
// 插入
DevicePointRulesDO devicePointRules = BeanUtils.toBean(createReqVO, DevicePointRulesDO.class);
devicePointRulesMapper.insert(devicePointRules);
mqttRuntimeCache.refreshDevice(devicePointRules.getDeviceId());
// 返回
return devicePointRules.getId();
}
@ -44,6 +48,7 @@ public class DevicePointRulesServiceImpl implements DevicePointRulesService {
public void updateDevicePointRules(DevicePointRulesSaveReqVO updateReqVO) {
// 校验存在
validateDevicePointRulesExists(updateReqVO.getId());
DevicePointRulesDO oldRule = devicePointRulesMapper.selectById(updateReqVO.getId());
// 更新
DevicePointRulesDO updateObj = BeanUtils.toBean(updateReqVO, DevicePointRulesDO.class);
// if (!updateReqVO.getPointRulesVOList().isEmpty()){
@ -53,14 +58,18 @@ public class DevicePointRulesServiceImpl implements DevicePointRulesService {
// }
devicePointRulesMapper.updateById(updateObj);
mqttRuntimeCache.refreshDevice(oldRule == null ? null : oldRule.getDeviceId());
mqttRuntimeCache.refreshDevice(updateObj.getDeviceId());
}
@Override
public void deleteDevicePointRules(Long id) {
// 校验存在
validateDevicePointRulesExists(id);
DevicePointRulesDO oldRule = devicePointRulesMapper.selectById(id);
// 删除
devicePointRulesMapper.deleteById(id);
mqttRuntimeCache.refreshDevice(oldRule == null ? null : oldRule.getDeviceId());
}
private void validateDevicePointRulesExists(Long id) {
@ -85,4 +94,4 @@ public class DevicePointRulesServiceImpl implements DevicePointRulesService {
.eq(DevicePointRulesDO::getDeviceId,id).eq(DevicePointRulesDO::getIdentifier,"ALARM").orderByDesc(DevicePointRulesDO::getDeviceId));
}
}
}

@ -30,7 +30,6 @@ import org.springframework.validation.annotation.Validated;
import javax.annotation.Resource;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.text.DecimalFormat;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
@ -53,6 +52,9 @@ import static cn.iocoder.yudao.module.mes.enums.ErrorCodeConstants.*;
@Slf4j
public class EnergyDeviceServiceImpl implements EnergyDeviceService {
private static final int POINT_VALUE_SCALE = 6;
private static final int ENERGY_VALUE_SCALE = 2;
@Resource
private EnergyDeviceMapper energyDeviceMapper;
@Resource
@ -1155,7 +1157,13 @@ public class EnergyDeviceServiceImpl implements EnergyDeviceService {
if (value == null) {
return null;
}
return ratio != null ? value * ratio : value;
BigDecimal result = BigDecimal.valueOf(value);
if (ratio != null) {
result = result.multiply(BigDecimal.valueOf(ratio));
}
return result.setScale(POINT_VALUE_SCALE, RoundingMode.HALF_UP)
.stripTrailingZeros()
.doubleValue();
}
/** 根据规则计算总值 */
@ -1317,8 +1325,11 @@ public class EnergyDeviceServiceImpl implements EnergyDeviceService {
/** 格式化 double */
private String formatDouble(Double value) {
if (value == null) return "0.0";
return new DecimalFormat("#.##").format(value);
if (value == null) return "0";
BigDecimal result = BigDecimal.valueOf(value)
.setScale(ENERGY_VALUE_SCALE, RoundingMode.HALF_UP)
.stripTrailingZeros();
return result.scale() < 0 ? result.setScale(0).toPlainString() : result.toPlainString();
}
/** 批量获取点位信息(本地数据库) */

Loading…
Cancel
Save