diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/config/DefaultEmqConfig.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/config/DefaultEmqConfig.java index 0bbdc758e..e2682a640 100644 --- a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/config/DefaultEmqConfig.java +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/config/DefaultEmqConfig.java @@ -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)); diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/config/MqttClientShutdown.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/config/MqttClientShutdown.java new file mode 100644 index 000000000..ff9de8b4f --- /dev/null +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/config/MqttClientShutdown.java @@ -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()); + } + } + +} diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/consumer/IotMqttRuntimeCache.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/consumer/IotMqttRuntimeCache.java new file mode 100644 index 000000000..1fd3f241a --- /dev/null +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/consumer/IotMqttRuntimeCache.java @@ -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> deviceByTopic = new ConcurrentHashMap<>(); + private final ConcurrentMap> pointsByDeviceId = new ConcurrentHashMap<>(); + private final ConcurrentMap>> pointRulesByDeviceId = new ConcurrentHashMap<>(); + private final ConcurrentMap> countRuleByDeviceId = new ConcurrentHashMap<>(); + + public DeviceDO getDeviceByTopic(String topic) { + if (StringUtils.isBlank(topic)) { + return null; + } + return deviceByTopic.computeIfAbsent(topic, this::loadDeviceByTopic).orElse(null); + } + + public List getDevicePoints(Long deviceId) { + if (deviceId == null) { + return Collections.emptyList(); + } + return pointsByDeviceId.computeIfAbsent(deviceId, this::loadDevicePoints); + } + + public List 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 loadDeviceByTopic(String topic) { + return Optional.ofNullable(deviceMapper.selectOne(Wrappers.lambdaQuery() + .eq(DeviceDO::getTopic, topic) + .last("LIMIT 1"))); + } + + private List loadDevicePoints(Long deviceId) { + List points = deviceContactModelMapper.selectList(Wrappers.lambdaQuery() + .eq(DeviceContactModelDO::getDeviceId, deviceId)); + if (CollectionUtils.isEmpty(points)) { + return Collections.emptyList(); + } + return Collections.unmodifiableList(new ArrayList<>(points)); + } + + private Map> loadPointRulesByAttributeCode(Long deviceId) { + List rules = devicePointRulesMapper.selectList(Wrappers.lambdaQuery() + .eq(DevicePointRulesDO::getDeviceId, deviceId) + .orderByDesc(DevicePointRulesDO::getCreateTime)); + if (CollectionUtils.isEmpty(rules)) { + return Collections.emptyMap(); + } + List parsedRules = new ArrayList<>(); + for (DevicePointRulesDO rule : rules) { + if (StringUtils.isBlank(rule.getFieldRule())) { + continue; + } + try { + List 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 loadLatestCountRule(Long deviceId) { + return Optional.ofNullable(devicePointRulesMapper.selectOne(Wrappers.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; + } + +} diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/consumer/MqttDataHandler.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/consumer/MqttDataHandler.java index 389418019..048f301f3 100644 --- a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/consumer/MqttDataHandler.java +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/consumer/MqttDataHandler.java @@ -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 { @Resource private IotOrganizationService organizationService; @Resource - @Lazy - private DeviceMapper deviceMapper; - @Resource private AsyncService asyncService; @Resource @@ -71,17 +64,14 @@ public class MqttDataHandler extends SuperConsumer { @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 { Map.class ); - DeviceDO deviceDO = deviceMapper.selectOne(Wrappers.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 { Long deviceId = device.getId(); // 1. 查询点位配置 - List points = getDevicePoints(deviceId); + List points = mqttRuntimeCache.getDevicePoints(deviceId); if (CollectionUtils.isEmpty(points)) { @@ -258,7 +248,7 @@ public class MqttDataHandler extends SuperConsumer { - 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 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 { 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 { // handleCapacityFormula private void handleCapacityFormula(DeviceDO device, Map varListMap) { - DevicePointRulesDO formulaRule = devicePointRulesMapper.selectOne( - Wrappers.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 { - private DevicePointRulesDO getDevicePointRules(Long deviceId) { - - List list = - devicePointRulesMapper.selectList( - Wrappers.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 getDevicePoints(Long deviceId) { - LambdaQueryWrapper query = new LambdaQueryWrapper<>(); - query.eq(DeviceContactModelDO::getDeviceId, deviceId); - return deviceContactModelMapper.selectList(query); - } - - /** * 处理OPC值 */ @@ -423,7 +368,7 @@ public class MqttDataHandler extends SuperConsumer { 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 { } } - 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 devicePointRulesDOList = devicePointRulesMapper.selectList( - Wrappers.lambdaQuery() - .eq(DevicePointRulesDO::getDeviceId, device.getId()).orderByDesc(DevicePointRulesDO::getCreateTime)); - - if (CollectionUtils.isEmpty(devicePointRulesDOList)) { + List 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 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 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 { 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"); diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/utils/ThreadUtils.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/utils/ThreadUtils.java index 39fe9daa7..6f18ab485 100644 --- a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/utils/ThreadUtils.java +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/framework/mqtt/utils/ThreadUtils.java @@ -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() + ); } diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/device/DeviceServiceImpl.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/device/DeviceServiceImpl.java index 43b2a7a98..cccf137d1 100644 --- a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/device/DeviceServiceImpl.java +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/device/DeviceServiceImpl.java @@ -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(); diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/device/TDengineService.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/device/TDengineService.java index 52e8e1cc4..8c5409b1f 100644 --- a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/device/TDengineService.java +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/device/TDengineService.java @@ -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> 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 existingColumns = getTDColumnNames(tableName); + Set 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 loadTDColumns(Long deviceId) { + return tdColumnCache.computeIfAbsent(deviceId, key -> ConcurrentHashMap.newKeySet()).isEmpty() + ? reloadTDColumns(deviceId) + : tdColumnCache.get(deviceId); + } + + private Set reloadTDColumns(Long deviceId) { + Set 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 cachedColumns = loadTDColumns(deviceId); + if (cachedColumns.contains(normalizedColumnName)) { + return true; + } + + Set 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) { diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/devicecontactmodel/DeviceContactModelServiceImpl.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/devicecontactmodel/DeviceContactModelServiceImpl.java index ea238c446..1b9483dcc 100644 --- a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/devicecontactmodel/DeviceContactModelServiceImpl.java +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/devicecontactmodel/DeviceContactModelServiceImpl.java @@ -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); } /** diff --git a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/devicepointrules/DevicePointRulesServiceImpl.java b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/devicepointrules/DevicePointRulesServiceImpl.java index fcfe23e23..e9c64bd51 100644 --- a/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/devicepointrules/DevicePointRulesServiceImpl.java +++ b/yudao-module-iot/yudao-module-iot-biz/src/main/java/cn/iocoder/yudao/module/iot/service/devicepointrules/DevicePointRulesServiceImpl.java @@ -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)); } -} \ No newline at end of file +} diff --git a/yudao-module-mes/yudao-module-mes-biz/src/main/java/cn/iocoder/yudao/module/mes/service/energydevice/EnergyDeviceServiceImpl.java b/yudao-module-mes/yudao-module-mes-biz/src/main/java/cn/iocoder/yudao/module/mes/service/energydevice/EnergyDeviceServiceImpl.java index 61dafd327..84c1f51a1 100644 --- a/yudao-module-mes/yudao-module-mes-biz/src/main/java/cn/iocoder/yudao/module/mes/service/energydevice/EnergyDeviceServiceImpl.java +++ b/yudao-module-mes/yudao-module-mes-biz/src/main/java/cn/iocoder/yudao/module/mes/service/energydevice/EnergyDeviceServiceImpl.java @@ -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(); } /** 批量获取点位信息(本地数据库) */